Skip to main content

ahri_tre_app/
configured_datastore_create.rs

1//! Privileged Datastore-creation workflow above configuration, Secret, and adapter capabilities.
2
3use ahri_tre_lake::{
4    AwsWorkloadIdentitySource, DatasetCsvSeed, DuckLakeAdapter, DuckLakeCredentialKind,
5    DuckLakeOpenedCatalog, FilesystemCatalogTls, FilesystemLakeSessionConfig,
6    ObjectLakeSessionConfig, ObjectStorageAuthentication, ObjectStorageLocation, ObjectStorageTls,
7    S3UrlStyle, ScratchAttempt, ScratchAttemptId, TrustedScratch,
8};
9use ahri_tre_pgmeta::{
10    PgDatastoreCreationConnection, PgDatastoreCreationSpec, PgDatastoreCreationStart,
11    PgDatastoreCredentialRotationError, PgDatastoreValidationSpec, PgPendingDatastoreCreation,
12    apply_datastore_credential_rotation, begin_datastore_creation, find_existing_datastore_binding,
13    inspect_datastore_identity_bindings, mark_datastore_identity_failed,
14    mark_datastore_identity_ready, read_datastore_identity_bindings,
15    rollback_datastore_credential_rotation, validate_existing_datastore_postgresql,
16};
17use ahri_tre_runtime::{
18    AdministrationFailureCode, AdministrationPeerIdentity, AdministrationRequest,
19    AdministrationRequestHandler, AdministrationResponse, AzureManagedIdentity,
20    ConfiguredAwsWorkloadIdentitySource, ConfiguredS3UrlStyle, ConfiguredSecret,
21    DatastoreConfigurationId, DatastoreCreationPlan, LakeStoragePlan, ManagedSecretOperationError,
22    PostgresqlBindingInspectionTarget, ResolvedOperationSecret, SessionTls,
23    StorageAuthenticationPlan, StorageTlsPlan, TrustedRuntimeState,
24};
25use ahri_tre_secrets::internal::{
26    ManagedSecretReferenceAuthorities, ManagedSecretReferenceAuthority,
27    ManagedSecretReferenceInspectionError, ManagedSecretReferenceStatus, ManagedSecretRemovalProof,
28    ManagedSecretRemovalProofErrorKind, StagedManagedSecret,
29};
30use ahri_tre_secrets::{
31    ManagedSecretErrorKind, ManagedSecretMetadata, ManagedSecretReference, SecretMaterial,
32};
33use ahri_tre_types::{
34    DatastoreCreationFailureCode, DatastoreIdentityBinding, DatastoreLakeLocation,
35    DatastoreLifecycleState, EncryptionMode, LakeCatalogCredentialMode, NcName,
36    NewDatastoreIdentityBinding, StudyId,
37};
38use chrono::{DateTime, Utc};
39use uuid::Uuid;
40
41const DEMO_DATASTORE: &str = "demo";
42const DEMO_STUDY_ID: &str = "00000000-0000-4000-8000-000000000015";
43
44/// Safe metadata for the binding-owned credential created during the saga.
45#[derive(Debug, Clone, PartialEq, Eq)]
46pub struct CreatedCatalogCredential {
47    version: u64,
48    created_at: DateTime<Utc>,
49}
50
51impl CreatedCatalogCredential {
52    pub fn new(version: u64, created_at: DateTime<Utc>) -> Self {
53        Self {
54            version,
55            created_at,
56        }
57    }
58
59    pub fn version(&self) -> u64 {
60        self.version
61    }
62
63    pub fn created_at(&self) -> DateTime<Utc> {
64        self.created_at
65    }
66}
67
68impl From<ManagedSecretMetadata> for CreatedCatalogCredential {
69    fn from(metadata: ManagedSecretMetadata) -> Self {
70        Self::new(metadata.version().get(), metadata.created_at())
71    }
72}
73/// Protected stage carried only by the Datastore credential owner workflow.
74pub struct StagedCatalogCredential {
75    inner: StagedManagedSecret,
76}
77
78impl StagedCatalogCredential {
79    pub fn version(&self) -> u64 {
80        self.inner.metadata().version().get()
81    }
82
83    pub fn material(&self) -> &SecretMaterial {
84        self.inner.material()
85    }
86}
87
88impl std::fmt::Debug for StagedCatalogCredential {
89    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
90        formatter
91            .debug_struct("StagedCatalogCredential")
92            .field("metadata", self.inner.metadata())
93            .finish()
94    }
95}
96
97/// One storage Secret paired with the exact configured reference that authorized it.
98pub struct ResolvedCreationStorageSecret {
99    configured: ConfiguredSecret,
100    resolved: ResolvedOperationSecret,
101}
102
103impl ResolvedCreationStorageSecret {
104    pub fn new(configured: ConfiguredSecret, resolved: ResolvedOperationSecret) -> Self {
105        Self {
106            configured,
107            resolved,
108        }
109    }
110
111    pub fn configured(&self) -> &ConfiguredSecret {
112        &self.configured
113    }
114
115    pub fn resolved(&self) -> &ResolvedOperationSecret {
116        &self.resolved
117    }
118}
119
120impl std::fmt::Debug for ResolvedCreationStorageSecret {
121    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
122        formatter
123            .debug_struct("ResolvedCreationStorageSecret")
124            .field("configured", &"<configured secret reference>")
125            .field("resolved", &self.resolved)
126            .finish()
127    }
128}
129
130/// Trusted authority retained by the handler. Clients never receive these capabilities.
131pub trait DatastoreCreationAuthority: Send + Sync {
132    fn deployment_id(&self) -> Result<Uuid, AdministrationFailureCode> {
133        Err(AdministrationFailureCode::ConfigurationUnavailable)
134    }
135
136    fn plan(
137        &self,
138        datastore: &DatastoreConfigurationId,
139    ) -> Result<DatastoreCreationPlan, AdministrationFailureCode>;
140
141    fn resolve_administrator(
142        &self,
143        plan: &DatastoreCreationPlan,
144    ) -> Result<ResolvedOperationSecret, AdministrationFailureCode>;
145
146    fn plan_for_existing(
147        &self,
148        datastore: &DatastoreConfigurationId,
149        datastore_id: uuid::Uuid,
150    ) -> Result<DatastoreCreationPlan, AdministrationFailureCode>;
151
152    fn resolve_existing_catalog_credential(
153        &self,
154        plan: &DatastoreCreationPlan,
155        binding: &DatastoreIdentityBinding,
156    ) -> Result<ResolvedOperationSecret, AdministrationFailureCode>;
157
158    fn resolve_storage(
159        &self,
160        plan: &DatastoreCreationPlan,
161    ) -> Result<Vec<ResolvedCreationStorageSecret>, AdministrationFailureCode>;
162
163    fn create_catalog_credential(
164        &self,
165        plan: &DatastoreCreationPlan,
166        material: &SecretMaterial,
167    ) -> Result<CreatedCatalogCredential, AdministrationFailureCode>;
168
169    fn verify_catalog_credential(
170        &self,
171        plan: &DatastoreCreationPlan,
172        material: &SecretMaterial,
173        credential: &CreatedCatalogCredential,
174    ) -> Result<(), AdministrationFailureCode>;
175
176    fn stage_catalog_credential_replacement(
177        &self,
178        _plan: &DatastoreCreationPlan,
179        _binding: &DatastoreIdentityBinding,
180        _material: &SecretMaterial,
181    ) -> Result<StagedCatalogCredential, AdministrationFailureCode> {
182        Err(AdministrationFailureCode::CredentialSetupFailed)
183    }
184
185    fn activate_catalog_credential_replacement(
186        &self,
187        _plan: &DatastoreCreationPlan,
188        _staged: &StagedCatalogCredential,
189    ) -> Result<CreatedCatalogCredential, AdministrationFailureCode> {
190        Err(AdministrationFailureCode::CredentialSetupFailed)
191    }
192
193    fn discard_catalog_credential_replacement(
194        &self,
195        _plan: &DatastoreCreationPlan,
196        _staged: &StagedCatalogCredential,
197    ) -> Result<(), AdministrationFailureCode> {
198        Err(AdministrationFailureCode::CredentialSetupFailed)
199    }
200
201    fn retain_datastore_binding_reference(
202        &self,
203        _plan: &DatastoreCreationPlan,
204        _reference: &ManagedSecretReference,
205    ) -> Result<(), AdministrationFailureCode> {
206        Err(AdministrationFailureCode::CredentialSetupFailed)
207    }
208
209    fn reconcilable_datastore_binding_reference_ids(
210        &self,
211        _reference: &ManagedSecretReference,
212    ) -> Result<Vec<Uuid>, ManagedSecretReferenceInspectionError> {
213        Ok(Vec::new())
214    }
215
216    fn release_absent_datastore_binding_reference(
217        &self,
218        _datastore_id: Uuid,
219        _reference: &ManagedSecretReference,
220    ) -> Result<(), ManagedSecretReferenceInspectionError> {
221        Err(ManagedSecretReferenceInspectionError::unavailable())
222    }
223
224    fn postgresql_binding_inspection_targets(
225        &self,
226    ) -> Result<Vec<PostgresqlBindingInspectionTarget>, AdministrationFailureCode> {
227        Err(AdministrationFailureCode::ConfigurationUnavailable)
228    }
229
230    fn resolve_postgresql_binding_inspection_administrator(
231        &self,
232        _target: &PostgresqlBindingInspectionTarget,
233    ) -> Result<ResolvedOperationSecret, AdministrationFailureCode> {
234        Err(AdministrationFailureCode::CredentialSetupFailed)
235    }
236
237    fn configured_datastore_ids(
238        &self,
239    ) -> Result<Vec<DatastoreConfigurationId>, AdministrationFailureCode> {
240        Err(AdministrationFailureCode::ConfigurationUnavailable)
241    }
242
243    fn application_reference_status(
244        &self,
245        _reference: &ManagedSecretReference,
246    ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
247        Err(ManagedSecretReferenceInspectionError::unavailable())
248    }
249
250    fn session_reference_status(
251        &self,
252        _reference: &ManagedSecretReference,
253    ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
254        Err(ManagedSecretReferenceInspectionError::unavailable())
255    }
256
257    fn datastore_reference_status(
258        &self,
259        _reference: &ManagedSecretReference,
260    ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
261        Err(ManagedSecretReferenceInspectionError::unavailable())
262    }
263
264    fn authentication_reference_status(
265        &self,
266        _reference: &ManagedSecretReference,
267    ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
268        Err(ManagedSecretReferenceInspectionError::unavailable())
269    }
270
271    fn remove_managed_secret(
272        &self,
273        _reference: &ManagedSecretReference,
274        _proof: ManagedSecretRemovalProof,
275    ) -> Result<ManagedSecretMetadata, AdministrationFailureCode> {
276        Err(AdministrationFailureCode::CredentialSetupFailed)
277    }
278}
279
280impl DatastoreCreationAuthority for TrustedRuntimeState {
281    fn deployment_id(&self) -> Result<Uuid, AdministrationFailureCode> {
282        Ok(TrustedRuntimeState::deployment_id(self))
283    }
284
285    fn plan(
286        &self,
287        datastore: &DatastoreConfigurationId,
288    ) -> Result<DatastoreCreationPlan, AdministrationFailureCode> {
289        self.datastore_creation_plan(datastore)
290            .map_err(|_| AdministrationFailureCode::ConfigurationUnavailable)
291    }
292
293    fn resolve_administrator(
294        &self,
295        plan: &DatastoreCreationPlan,
296    ) -> Result<ResolvedOperationSecret, AdministrationFailureCode> {
297        self.resolve_datastore_administrator(plan)
298            .map_err(|_| AdministrationFailureCode::CredentialSetupFailed)
299    }
300
301    fn plan_for_existing(
302        &self,
303        datastore: &DatastoreConfigurationId,
304        datastore_id: uuid::Uuid,
305    ) -> Result<DatastoreCreationPlan, AdministrationFailureCode> {
306        self.datastore_creation_plan_with_id(datastore, datastore_id)
307            .map_err(|_| AdministrationFailureCode::ConfigurationUnavailable)
308    }
309
310    fn resolve_existing_catalog_credential(
311        &self,
312        plan: &DatastoreCreationPlan,
313        binding: &DatastoreIdentityBinding,
314    ) -> Result<ResolvedOperationSecret, AdministrationFailureCode> {
315        TrustedRuntimeState::resolve_existing_datastore_catalog_credential(self, plan, binding)
316            .map_err(|_| AdministrationFailureCode::CredentialSetupFailed)
317    }
318
319    fn resolve_storage(
320        &self,
321        plan: &DatastoreCreationPlan,
322    ) -> Result<Vec<ResolvedCreationStorageSecret>, AdministrationFailureCode> {
323        plan.storage_secrets()
324            .into_iter()
325            .map(|configured| {
326                self.resolve_datastore_creation_storage(plan, configured)
327                    .map(|resolved| {
328                        ResolvedCreationStorageSecret::new(configured.clone(), resolved)
329                    })
330                    .map_err(|_| AdministrationFailureCode::CredentialSetupFailed)
331            })
332            .collect()
333    }
334
335    fn create_catalog_credential(
336        &self,
337        plan: &DatastoreCreationPlan,
338        material: &SecretMaterial,
339    ) -> Result<CreatedCatalogCredential, AdministrationFailureCode> {
340        TrustedRuntimeState::create_datastore_catalog_credential(self, plan, material)
341            .map(Into::into)
342            .map_err(|_| AdministrationFailureCode::CredentialSetupFailed)
343    }
344
345    fn verify_catalog_credential(
346        &self,
347        plan: &DatastoreCreationPlan,
348        material: &SecretMaterial,
349        credential: &CreatedCatalogCredential,
350    ) -> Result<(), AdministrationFailureCode> {
351        TrustedRuntimeState::verify_datastore_catalog_credential(
352            self,
353            plan,
354            material,
355            credential.version(),
356        )
357        .map_err(|_| AdministrationFailureCode::CredentialSetupFailed)
358    }
359
360    fn stage_catalog_credential_replacement(
361        &self,
362        plan: &DatastoreCreationPlan,
363        binding: &DatastoreIdentityBinding,
364        material: &SecretMaterial,
365    ) -> Result<StagedCatalogCredential, AdministrationFailureCode> {
366        self.stage_datastore_catalog_credential_replacement(plan, binding, material)
367            .map(|inner| StagedCatalogCredential { inner })
368            .map_err(|_| AdministrationFailureCode::CredentialSetupFailed)
369    }
370
371    fn activate_catalog_credential_replacement(
372        &self,
373        plan: &DatastoreCreationPlan,
374        staged: &StagedCatalogCredential,
375    ) -> Result<CreatedCatalogCredential, AdministrationFailureCode> {
376        self.activate_datastore_catalog_credential_replacement(plan, &staged.inner)
377            .map(Into::into)
378            .map_err(|error| match error {
379                ManagedSecretOperationError::Resolve(error)
380                    if error.kind() == ManagedSecretErrorKind::CommitUncertain =>
381                {
382                    AdministrationFailureCode::CredentialStateUncertain
383                }
384                _ => AdministrationFailureCode::CredentialSetupFailed,
385            })
386    }
387
388    fn discard_catalog_credential_replacement(
389        &self,
390        plan: &DatastoreCreationPlan,
391        staged: &StagedCatalogCredential,
392    ) -> Result<(), AdministrationFailureCode> {
393        self.discard_datastore_catalog_credential_replacement(plan, &staged.inner)
394            .map_err(|_| AdministrationFailureCode::CredentialSetupFailed)
395    }
396
397    fn retain_datastore_binding_reference(
398        &self,
399        plan: &DatastoreCreationPlan,
400        reference: &ManagedSecretReference,
401    ) -> Result<(), AdministrationFailureCode> {
402        TrustedRuntimeState::retain_datastore_binding_reference(self, plan, reference)
403            .map_err(|_| AdministrationFailureCode::CredentialSetupFailed)
404    }
405
406    fn reconcilable_datastore_binding_reference_ids(
407        &self,
408        reference: &ManagedSecretReference,
409    ) -> Result<Vec<Uuid>, ManagedSecretReferenceInspectionError> {
410        TrustedRuntimeState::reconcilable_datastore_binding_reference_ids(self, reference)
411            .map_err(|_| ManagedSecretReferenceInspectionError::unavailable())
412    }
413
414    fn release_absent_datastore_binding_reference(
415        &self,
416        datastore_id: Uuid,
417        reference: &ManagedSecretReference,
418    ) -> Result<(), ManagedSecretReferenceInspectionError> {
419        TrustedRuntimeState::release_absent_datastore_binding_reference(
420            self,
421            datastore_id,
422            reference,
423        )
424        .map_err(|_| ManagedSecretReferenceInspectionError::unavailable())
425    }
426
427    fn postgresql_binding_inspection_targets(
428        &self,
429    ) -> Result<Vec<PostgresqlBindingInspectionTarget>, AdministrationFailureCode> {
430        TrustedRuntimeState::postgresql_binding_inspection_targets(self)
431            .map_err(|_| AdministrationFailureCode::ConfigurationUnavailable)
432    }
433
434    fn resolve_postgresql_binding_inspection_administrator(
435        &self,
436        target: &PostgresqlBindingInspectionTarget,
437    ) -> Result<ResolvedOperationSecret, AdministrationFailureCode> {
438        TrustedRuntimeState::resolve_postgresql_binding_inspection_administrator(self, target)
439            .map_err(|_| AdministrationFailureCode::CredentialSetupFailed)
440    }
441
442    fn configured_datastore_ids(
443        &self,
444    ) -> Result<Vec<DatastoreConfigurationId>, AdministrationFailureCode> {
445        Ok(TrustedRuntimeState::configured_datastore_ids(self))
446    }
447
448    fn application_reference_status(
449        &self,
450        reference: &ManagedSecretReference,
451    ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
452        TrustedRuntimeState::application_reference_status(self, reference)
453    }
454
455    fn session_reference_status(
456        &self,
457        reference: &ManagedSecretReference,
458    ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
459        TrustedRuntimeState::session_reference_status(self, reference)
460    }
461
462    fn datastore_reference_status(
463        &self,
464        reference: &ManagedSecretReference,
465    ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
466        TrustedRuntimeState::datastore_reference_status(self, reference)
467    }
468
469    fn authentication_reference_status(
470        &self,
471        reference: &ManagedSecretReference,
472    ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
473        TrustedRuntimeState::authentication_reference_status(self, reference)
474    }
475
476    fn remove_managed_secret(
477        &self,
478        reference: &ManagedSecretReference,
479        proof: ManagedSecretRemovalProof,
480    ) -> Result<ManagedSecretMetadata, AdministrationFailureCode> {
481        TrustedRuntimeState::remove_managed_secret(self, reference, proof)
482            .map_err(|_| AdministrationFailureCode::CredentialSetupFailed)
483    }
484}
485
486impl<A> DatastoreCreationAuthority for std::sync::Arc<A>
487where
488    A: DatastoreCreationAuthority,
489{
490    fn deployment_id(&self) -> Result<Uuid, AdministrationFailureCode> {
491        (**self).deployment_id()
492    }
493
494    fn plan(
495        &self,
496        datastore: &DatastoreConfigurationId,
497    ) -> Result<DatastoreCreationPlan, AdministrationFailureCode> {
498        (**self).plan(datastore)
499    }
500
501    fn resolve_administrator(
502        &self,
503        plan: &DatastoreCreationPlan,
504    ) -> Result<ResolvedOperationSecret, AdministrationFailureCode> {
505        (**self).resolve_administrator(plan)
506    }
507
508    fn plan_for_existing(
509        &self,
510        datastore: &DatastoreConfigurationId,
511        datastore_id: uuid::Uuid,
512    ) -> Result<DatastoreCreationPlan, AdministrationFailureCode> {
513        (**self).plan_for_existing(datastore, datastore_id)
514    }
515
516    fn resolve_existing_catalog_credential(
517        &self,
518        plan: &DatastoreCreationPlan,
519        binding: &DatastoreIdentityBinding,
520    ) -> Result<ResolvedOperationSecret, AdministrationFailureCode> {
521        (**self).resolve_existing_catalog_credential(plan, binding)
522    }
523
524    fn resolve_storage(
525        &self,
526        plan: &DatastoreCreationPlan,
527    ) -> Result<Vec<ResolvedCreationStorageSecret>, AdministrationFailureCode> {
528        (**self).resolve_storage(plan)
529    }
530
531    fn create_catalog_credential(
532        &self,
533        plan: &DatastoreCreationPlan,
534        material: &SecretMaterial,
535    ) -> Result<CreatedCatalogCredential, AdministrationFailureCode> {
536        (**self).create_catalog_credential(plan, material)
537    }
538
539    fn verify_catalog_credential(
540        &self,
541        plan: &DatastoreCreationPlan,
542        material: &SecretMaterial,
543        credential: &CreatedCatalogCredential,
544    ) -> Result<(), AdministrationFailureCode> {
545        (**self).verify_catalog_credential(plan, material, credential)
546    }
547
548    fn stage_catalog_credential_replacement(
549        &self,
550        plan: &DatastoreCreationPlan,
551        binding: &DatastoreIdentityBinding,
552        material: &SecretMaterial,
553    ) -> Result<StagedCatalogCredential, AdministrationFailureCode> {
554        (**self).stage_catalog_credential_replacement(plan, binding, material)
555    }
556
557    fn activate_catalog_credential_replacement(
558        &self,
559        plan: &DatastoreCreationPlan,
560        staged: &StagedCatalogCredential,
561    ) -> Result<CreatedCatalogCredential, AdministrationFailureCode> {
562        (**self).activate_catalog_credential_replacement(plan, staged)
563    }
564
565    fn discard_catalog_credential_replacement(
566        &self,
567        plan: &DatastoreCreationPlan,
568        staged: &StagedCatalogCredential,
569    ) -> Result<(), AdministrationFailureCode> {
570        (**self).discard_catalog_credential_replacement(plan, staged)
571    }
572
573    fn retain_datastore_binding_reference(
574        &self,
575        plan: &DatastoreCreationPlan,
576        reference: &ManagedSecretReference,
577    ) -> Result<(), AdministrationFailureCode> {
578        (**self).retain_datastore_binding_reference(plan, reference)
579    }
580
581    fn reconcilable_datastore_binding_reference_ids(
582        &self,
583        reference: &ManagedSecretReference,
584    ) -> Result<Vec<Uuid>, ManagedSecretReferenceInspectionError> {
585        (**self).reconcilable_datastore_binding_reference_ids(reference)
586    }
587
588    fn release_absent_datastore_binding_reference(
589        &self,
590        datastore_id: Uuid,
591        reference: &ManagedSecretReference,
592    ) -> Result<(), ManagedSecretReferenceInspectionError> {
593        (**self).release_absent_datastore_binding_reference(datastore_id, reference)
594    }
595
596    fn postgresql_binding_inspection_targets(
597        &self,
598    ) -> Result<Vec<PostgresqlBindingInspectionTarget>, AdministrationFailureCode> {
599        (**self).postgresql_binding_inspection_targets()
600    }
601
602    fn resolve_postgresql_binding_inspection_administrator(
603        &self,
604        target: &PostgresqlBindingInspectionTarget,
605    ) -> Result<ResolvedOperationSecret, AdministrationFailureCode> {
606        (**self).resolve_postgresql_binding_inspection_administrator(target)
607    }
608
609    fn configured_datastore_ids(
610        &self,
611    ) -> Result<Vec<DatastoreConfigurationId>, AdministrationFailureCode> {
612        (**self).configured_datastore_ids()
613    }
614
615    fn application_reference_status(
616        &self,
617        reference: &ManagedSecretReference,
618    ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
619        (**self).application_reference_status(reference)
620    }
621
622    fn session_reference_status(
623        &self,
624        reference: &ManagedSecretReference,
625    ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
626        (**self).session_reference_status(reference)
627    }
628
629    fn datastore_reference_status(
630        &self,
631        reference: &ManagedSecretReference,
632    ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
633        (**self).datastore_reference_status(reference)
634    }
635
636    fn authentication_reference_status(
637        &self,
638        reference: &ManagedSecretReference,
639    ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
640        (**self).authentication_reference_status(reference)
641    }
642
643    fn remove_managed_secret(
644        &self,
645        reference: &ManagedSecretReference,
646        proof: ManagedSecretRemovalProof,
647    ) -> Result<ManagedSecretMetadata, AdministrationFailureCode> {
648        (**self).remove_managed_secret(reference, proof)
649    }
650}
651
652/// Adapter-owned phases of the cross-store creation saga.
653pub trait DatastoreCreationProvisioner: Send + Sync {
654    type Pending;
655
656    /// Creates PostgreSQL resources and persists the singleton `creating` binding.
657    fn begin(
658        &self,
659        plan: &DatastoreCreationPlan,
660        _administrator: &ResolvedOperationSecret,
661        catalog_credential: &SecretMaterial,
662    ) -> DatastoreCreationStart<Self::Pending>;
663
664    fn initialize_lake(
665        &self,
666        pending: &mut Self::Pending,
667        plan: &DatastoreCreationPlan,
668        catalog_credential: &SecretMaterial,
669        storage: &[ResolvedCreationStorageSecret],
670    ) -> Result<(), DatastoreCreationFailure>;
671
672    fn verify(
673        &self,
674        pending: &mut Self::Pending,
675        plan: &DatastoreCreationPlan,
676    ) -> Result<(), DatastoreCreationFailure>;
677
678    fn mark_ready(
679        &self,
680        pending: &mut Self::Pending,
681        credential: &CreatedCatalogCredential,
682    ) -> Result<(), DatastoreCreationFailure>;
683
684    fn mark_failed(
685        &self,
686        pending: &mut Self::Pending,
687        failure: &DatastoreCreationFailure,
688    ) -> Result<(), DatastoreCreationFailure>;
689
690    fn find_existing(
691        &self,
692        _plan: &DatastoreCreationPlan,
693        _administrator: &ResolvedOperationSecret,
694    ) -> Result<Option<DatastoreIdentityBinding>, DatastoreCreationFailure> {
695        Ok(None)
696    }
697
698    fn inspect_all_bindings(
699        &self,
700        _target: &PostgresqlBindingInspectionTarget,
701        _administrator: &ResolvedOperationSecret,
702    ) -> Result<Vec<DatastoreIdentityBinding>, DatastoreCreationFailure> {
703        Err(verification_failure(()))
704    }
705
706    fn verify_existing(
707        &self,
708        _plan: &DatastoreCreationPlan,
709        _binding: &DatastoreIdentityBinding,
710        _administrator: &ResolvedOperationSecret,
711        _catalog_credential: &ResolvedOperationSecret,
712        _storage: &[ResolvedCreationStorageSecret],
713    ) -> Result<(), DatastoreCreationFailure> {
714        Err(verification_failure(()))
715    }
716
717    fn plan_governance(
718        &self,
719        _plan: &DatastoreCreationPlan,
720        _binding: &DatastoreIdentityBinding,
721        _administrator: &ResolvedOperationSecret,
722        _catalog: &ResolvedOperationSecret,
723        _storage: &[ResolvedCreationStorageSecret],
724    ) -> Result<ahri_tre_pgmeta::governance::GovernanceUpgradePlan, DatastoreCreationFailure> {
725        Err(verification_failure(()))
726    }
727
728    fn upgrade_governance(
729        &self,
730        _plan: &DatastoreCreationPlan,
731        _binding: &DatastoreIdentityBinding,
732        _administrator: &ResolvedOperationSecret,
733        _catalog_credential: &ResolvedOperationSecret,
734        _storage: &[ResolvedCreationStorageSecret],
735        _batch: u32,
736    ) -> Result<GovernanceUpgradeProgress, DatastoreCreationFailure> {
737        Err(verification_failure(()))
738    }
739
740    fn seed_demo(
741        &self,
742        _plan: &DatastoreCreationPlan,
743        _binding: &DatastoreIdentityBinding,
744        _administrator: &ResolvedOperationSecret,
745        _catalog_credential: &ResolvedOperationSecret,
746        _storage: &[ResolvedCreationStorageSecret],
747    ) -> Result<(), DatastoreCreationFailure> {
748        Err(verification_failure(()))
749    }
750
751    fn apply_catalog_credential_rotation(
752        &self,
753        _plan: &DatastoreCreationPlan,
754        _binding: &DatastoreIdentityBinding,
755        _administrator: &ResolvedOperationSecret,
756        _prior: &ResolvedOperationSecret,
757        _staged: &StagedCatalogCredential,
758        _rotated_at: DateTime<Utc>,
759    ) -> Result<DatastoreIdentityBinding, DatastoreCreationFailure> {
760        Err(verification_failure(()))
761    }
762
763    fn rollback_catalog_credential_rotation(
764        &self,
765        _plan: &DatastoreCreationPlan,
766        _binding: &DatastoreIdentityBinding,
767        _administrator: &ResolvedOperationSecret,
768        _prior: &ResolvedOperationSecret,
769    ) -> Result<(), DatastoreCreationFailure> {
770        Err(verification_failure(()))
771    }
772}
773
774/// Production PostgreSQL and DuckLake adapter implementation of the creation saga.
775#[derive(Debug, Default, Clone, Copy)]
776pub struct ConfiguredDatastoreCreationProvisioner;
777
778pub struct PendingConfiguredDatastoreCreation {
779    postgres: PgPendingDatastoreCreation,
780    lake: Option<DuckLakeOpenedCatalog>,
781    scratch: Option<ScratchAttempt>,
782}
783
784impl std::fmt::Debug for PendingConfiguredDatastoreCreation {
785    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
786        formatter
787            .debug_struct("PendingConfiguredDatastoreCreation")
788            .field("postgres", &self.postgres)
789            .field("lake_initialized", &self.lake.is_some())
790            .finish_non_exhaustive()
791    }
792}
793
794#[derive(Debug)]
795pub enum DatastoreCreationStart<P> {
796    Started(P),
797    Conflict,
798    Failed(DatastoreCreationFailure),
799    FailedWithBinding {
800        pending: P,
801        failure: DatastoreCreationFailure,
802    },
803}
804
805#[derive(Debug, Clone, PartialEq, Eq)]
806pub struct DatastoreCreationFailure {
807    code: DatastoreCreationFailureCode,
808    safe_summary: &'static str,
809}
810
811impl DatastoreCreationFailure {
812    pub fn new(code: DatastoreCreationFailureCode, safe_summary: &'static str) -> Self {
813        debug_assert!(safe_summary.len() <= 512);
814        Self { code, safe_summary }
815    }
816
817    pub fn code(&self) -> DatastoreCreationFailureCode {
818        self.code
819    }
820
821    pub fn safe_summary(&self) -> &'static str {
822        self.safe_summary
823    }
824
825    fn administration_code(&self) -> AdministrationFailureCode {
826        match self.code {
827            DatastoreCreationFailureCode::PostgresqlSetupFailed => {
828                AdministrationFailureCode::PostgresqlSetupFailed
829            }
830            DatastoreCreationFailureCode::CredentialSetupFailed => {
831                AdministrationFailureCode::CredentialSetupFailed
832            }
833            DatastoreCreationFailureCode::LakeSetupFailed => {
834                AdministrationFailureCode::LakeSetupFailed
835            }
836            DatastoreCreationFailureCode::VerificationFailed => {
837                AdministrationFailureCode::VerificationFailed
838            }
839            DatastoreCreationFailureCode::InternalFailed => {
840                AdministrationFailureCode::InternalFailed
841            }
842            DatastoreCreationFailureCode::CredentialStateUncertain => {
843                AdministrationFailureCode::CredentialStateUncertain
844            }
845        }
846    }
847}
848
849impl DatastoreCreationProvisioner for ConfiguredDatastoreCreationProvisioner {
850    type Pending = PendingConfiguredDatastoreCreation;
851
852    fn begin(
853        &self,
854        plan: &DatastoreCreationPlan,
855        administrator: &ResolvedOperationSecret,
856        catalog_credential: &SecretMaterial,
857    ) -> DatastoreCreationStart<Self::Pending> {
858        let (lake_path, lake_location) = match creation_lake_location(plan) {
859            Ok(location) => location,
860            Err(()) => return DatastoreCreationStart::Failed(postgresql_failure()),
861        };
862        let binding = match NewDatastoreIdentityBinding::creating_managed(
863            plan.datastore_id(),
864            plan.postgresql().database(),
865            plan.catalog_database(),
866            plan.catalog_schema(),
867            lake_path,
868            plan.storage_policy_id(),
869            lake_location,
870            EncryptionMode::ServerManaged,
871            plan.catalog_role(),
872            plan.catalog_secret_reference().as_str(),
873            Some(env!("CARGO_PKG_VERSION").to_string()),
874        ) {
875            Ok(binding) => binding,
876            Err(_) => return DatastoreCreationStart::Failed(postgresql_failure()),
877        };
878        let connection = creation_postgres_connection(plan);
879        let Ok(governance) = ahri_tre_pgmeta::governance::GovernanceAppointment::from_identities(
880            plan.governance_custodians()
881                .iter()
882                .map(|a| (a.issuer(), a.subject())),
883        ) else {
884            return DatastoreCreationStart::Failed(verification_failure(()));
885        };
886        let spec = PgDatastoreCreationSpec::new(
887            plan.postgresql().database(),
888            plan.metadata_role(),
889            plan.catalog_database(),
890            plan.catalog_schema(),
891            plan.catalog_role(),
892            binding,
893            governance,
894        );
895        let started = begin_datastore_creation(
896            &connection,
897            administrator.material(),
898            catalog_credential,
899            spec,
900        );
901        match started {
902            PgDatastoreCreationStart::Started(postgres) => {
903                DatastoreCreationStart::Started(PendingConfiguredDatastoreCreation {
904                    postgres,
905                    lake: None,
906                    scratch: None,
907                })
908            }
909            PgDatastoreCreationStart::Conflict => DatastoreCreationStart::Conflict,
910            PgDatastoreCreationStart::FailedBeforeBinding => {
911                DatastoreCreationStart::Failed(postgresql_failure())
912            }
913            PgDatastoreCreationStart::RollbackFailed => {
914                DatastoreCreationStart::Failed(internal_failure())
915            }
916            PgDatastoreCreationStart::FailedAfterBinding(postgres) => {
917                DatastoreCreationStart::FailedWithBinding {
918                    pending: PendingConfiguredDatastoreCreation {
919                        postgres,
920                        lake: None,
921                        scratch: None,
922                    },
923                    failure: postgresql_failure(),
924                }
925            }
926        }
927    }
928
929    fn initialize_lake(
930        &self,
931        pending: &mut Self::Pending,
932        plan: &DatastoreCreationPlan,
933        catalog_credential: &SecretMaterial,
934        storage: &[ResolvedCreationStorageSecret],
935    ) -> Result<(), DatastoreCreationFailure> {
936        let catalog_tls = catalog_tls(plan.postgresql().tls());
937        let (opened, scratch) = match plan.lake().storage() {
938            LakeStoragePlan::Filesystem { base_path } => {
939                let config = FilesystemLakeSessionConfig::new(
940                    base_path,
941                    plan.datastore_prefix(),
942                    plan.postgresql().host(),
943                    plan.postgresql().port(),
944                    plan.catalog_database(),
945                    plan.catalog_schema(),
946                    plan.catalog_role(),
947                    catalog_tls,
948                    EncryptionMode::ServerManaged,
949                )
950                .map_err(lake_failure)?;
951                let scratch = create_scratch(plan, [config.canonical_data_path()])?;
952                let opened = DuckLakeAdapter::initialize_filesystem_catalog_with_scratch(
953                    &config,
954                    catalog_credential,
955                    &scratch,
956                )
957                .map_err(lake_failure)?;
958                (opened, scratch)
959            }
960            _ => {
961                let location = object_creation_location(plan).map_err(lake_failure)?;
962                let config = ObjectLakeSessionConfig::new(
963                    location,
964                    plan.postgresql().host(),
965                    plan.postgresql().port(),
966                    plan.catalog_database(),
967                    plan.catalog_schema(),
968                    plan.catalog_role(),
969                    catalog_tls,
970                    object_storage_tls(plan.lake().storage()),
971                    EncryptionMode::ServerManaged,
972                )
973                .map_err(lake_failure)?;
974                let scratch = create_scratch(plan, std::iter::empty::<String>())?;
975                let opened =
976                    initialize_object_lake(&config, catalog_credential, plan, storage, &scratch)?;
977                (opened, scratch)
978            }
979        };
980        pending.lake = Some(opened);
981        pending.scratch = Some(scratch);
982        Ok(())
983    }
984
985    fn verify(
986        &self,
987        pending: &mut Self::Pending,
988        plan: &DatastoreCreationPlan,
989    ) -> Result<(), DatastoreCreationFailure> {
990        let bindings = read_datastore_identity_bindings(pending.postgres.store_mut())
991            .map_err(verification_failure)?;
992        let [binding] = bindings.as_slice() else {
993            return Err(verification_failure(()));
994        };
995        let lake = pending
996            .lake
997            .as_ref()
998            .ok_or_else(|| verification_failure(()))?;
999        let (expected_path, expected_location) =
1000            creation_lake_location(plan).map_err(verification_failure)?;
1001        if binding.datastore_id != plan.datastore_id()
1002            || binding.datastore_name != plan.postgresql().database()
1003            || binding.ducklake_catalog_database != plan.catalog_database()
1004            || binding.ducklake_catalog_schema.as_deref() != Some(plan.catalog_schema())
1005            || binding.storage_policy_id.as_deref() != Some(plan.storage_policy_id())
1006            || binding.lake_location.as_ref() != Some(&expected_location)
1007            || binding.lake_catalog_role_name.as_deref() != Some(plan.catalog_role())
1008            || binding.managed_secret_ref.as_deref()
1009                != Some(plan.catalog_secret_reference().as_str())
1010            || binding.lake_catalog_credential_mode != LakeCatalogCredentialMode::ManagedLocal
1011            || binding.credential_version != 1
1012            || !binding.ducklake_encryption
1013            || binding.lifecycle_state != DatastoreLifecycleState::Creating
1014            || binding.binding_fingerprint != binding.computed_binding_fingerprint()
1015            || binding.lake_path != expected_path.as_str()
1016            || !equivalent_ducklake_data_paths(&binding.lake_path, &lake.health.data_path)
1017            || lake.attach_plan.data_path != expected_path.as_str()
1018            || lake.attach_plan.credential_kind != DuckLakeCredentialKind::SessionSecret
1019            || lake.attach_plan.configured_encryption_mode != EncryptionMode::ServerManaged
1020            || lake.health.detected_encryption_mode != EncryptionMode::ServerManaged
1021            || !lake.attach_plan.create_if_not_exists
1022            || lake.attach_plan.automatic_migration
1023        {
1024            return Err(verification_failure(()));
1025        }
1026        Ok(())
1027    }
1028
1029    fn mark_ready(
1030        &self,
1031        pending: &mut Self::Pending,
1032        credential: &CreatedCatalogCredential,
1033    ) -> Result<(), DatastoreCreationFailure> {
1034        if credential.version() != 1 {
1035            return Err(verification_failure(()));
1036        }
1037        let bindings = read_datastore_identity_bindings(pending.postgres.store_mut())
1038            .map_err(verification_failure)?;
1039        let [binding] = bindings.as_slice() else {
1040            return Err(verification_failure(()));
1041        };
1042        let ledger = ahri_tre_pgmeta::governance::ledger_binding(pending.postgres.store_mut())
1043            .map_err(verification_failure)?;
1044        let lake = pending
1045            .lake
1046            .as_mut()
1047            .ok_or_else(|| verification_failure(()))?;
1048        ahri_tre_lake::GovernanceLedger::provision(
1049            &mut lake.connection,
1050            ahri_tre_types::StudyId(ledger.study_id),
1051        )
1052        .map_err(verification_failure)?;
1053        crate::disclosure::reconcile_governance_evidence(
1054            pending.postgres.store_mut(),
1055            &mut lake.connection,
1056            1000,
1057        )
1058        .map_err(verification_failure)?;
1059        mark_datastore_identity_ready(
1060            pending.postgres.store_mut(),
1061            binding.datastore_id,
1062            credential.created_at(),
1063        )
1064        .map_err(verification_failure)?;
1065        Ok(())
1066    }
1067
1068    fn mark_failed(
1069        &self,
1070        pending: &mut Self::Pending,
1071        failure: &DatastoreCreationFailure,
1072    ) -> Result<(), DatastoreCreationFailure> {
1073        let bindings = read_datastore_identity_bindings(pending.postgres.store_mut())
1074            .map_err(|_| internal_failure())?;
1075        let [binding] = bindings.as_slice() else {
1076            return Err(internal_failure());
1077        };
1078        mark_datastore_identity_failed(
1079            pending.postgres.store_mut(),
1080            binding.datastore_id,
1081            failure.code(),
1082            Some(failure.safe_summary()),
1083        )
1084        .map_err(|_| internal_failure())?;
1085        Ok(())
1086    }
1087
1088    fn find_existing(
1089        &self,
1090        plan: &DatastoreCreationPlan,
1091        administrator: &ResolvedOperationSecret,
1092    ) -> Result<Option<DatastoreIdentityBinding>, DatastoreCreationFailure> {
1093        find_existing_datastore_binding(
1094            &creation_postgres_connection(plan),
1095            administrator.material(),
1096            plan.postgresql().database(),
1097        )
1098        .map_err(verification_failure)
1099    }
1100
1101    fn inspect_all_bindings(
1102        &self,
1103        target: &PostgresqlBindingInspectionTarget,
1104        administrator: &ResolvedOperationSecret,
1105    ) -> Result<Vec<DatastoreIdentityBinding>, DatastoreCreationFailure> {
1106        let connection = PgDatastoreCreationConnection::new(
1107            target.host(),
1108            target.port(),
1109            target.bootstrap_database(),
1110            target.administrator_username(),
1111            target.root_certificates().to_vec(),
1112        );
1113        inspect_datastore_identity_bindings(&connection, administrator.material())
1114            .map_err(|_| postgresql_failure())
1115    }
1116
1117    fn verify_existing(
1118        &self,
1119        plan: &DatastoreCreationPlan,
1120        binding: &DatastoreIdentityBinding,
1121        administrator: &ResolvedOperationSecret,
1122        catalog_credential: &ResolvedOperationSecret,
1123        storage: &[ResolvedCreationStorageSecret],
1124    ) -> Result<(), DatastoreCreationFailure> {
1125        let (expected_path, expected_location) =
1126            creation_lake_location(plan).map_err(verification_failure)?;
1127        if binding.datastore_id != plan.datastore_id()
1128            || binding.datastore_name != plan.postgresql().database()
1129            || binding.ducklake_catalog_database != plan.catalog_database()
1130            || binding.ducklake_catalog_schema.as_deref() != Some(plan.catalog_schema())
1131            || binding.storage_policy_id.as_deref() != Some(plan.storage_policy_id())
1132            || binding.lake_location.as_ref() != Some(&expected_location)
1133            || binding.lake_catalog_role_name.as_deref() != Some(plan.catalog_role())
1134            || binding.managed_secret_ref.as_deref()
1135                != Some(plan.catalog_secret_reference().as_str())
1136            || binding.lake_catalog_credential_mode != LakeCatalogCredentialMode::ManagedLocal
1137            || binding.credential_version
1138                != i32::try_from(catalog_credential.version().managed_version().unwrap_or(0))
1139                    .map_err(verification_failure)?
1140            || !binding.ducklake_encryption
1141            || binding.lifecycle_state != DatastoreLifecycleState::Ready
1142            || binding.failure_code.is_some()
1143            || binding.binding_fingerprint != binding.computed_binding_fingerprint()
1144            || binding.lake_path != expected_path
1145        {
1146            return Err(verification_failure(()));
1147        }
1148        let postgres_spec = PgDatastoreValidationSpec::new(
1149            plan.postgresql().database(),
1150            plan.metadata_role(),
1151            plan.catalog_database(),
1152            plan.catalog_schema(),
1153            plan.catalog_role(),
1154        );
1155        validate_existing_datastore_postgresql(
1156            &creation_postgres_connection(plan),
1157            administrator.material(),
1158            &postgres_spec,
1159        )
1160        .map_err(verification_failure)?;
1161
1162        let scratch = create_scratch(plan, std::iter::once(expected_path.clone()))?;
1163        let opened = open_existing_lake(plan, catalog_credential.material(), storage, &scratch)?;
1164        if !equivalent_ducklake_data_paths(&opened.health.data_path, &expected_path)
1165            || opened.attach_plan.data_path != expected_path
1166            || opened.attach_plan.create_if_not_exists
1167            || opened.attach_plan.configured_encryption_mode != EncryptionMode::ServerManaged
1168            || opened.health.detected_encryption_mode != EncryptionMode::ServerManaged
1169        {
1170            return Err(verification_failure(()));
1171        }
1172        Ok(())
1173    }
1174
1175    fn plan_governance(
1176        &self,
1177        plan: &DatastoreCreationPlan,
1178        binding: &DatastoreIdentityBinding,
1179        administrator: &ResolvedOperationSecret,
1180        catalog: &ResolvedOperationSecret,
1181        storage: &[ResolvedCreationStorageSecret],
1182    ) -> Result<ahri_tre_pgmeta::governance::GovernanceUpgradePlan, DatastoreCreationFailure> {
1183        self.verify_existing(plan, binding, administrator, catalog, storage)?;
1184        creation_postgres_connection(plan)
1185            .dataset_maintenance(plan.postgresql().database(), administrator.material())
1186            .map_err(verification_failure)?
1187            .governance()
1188            .upgrade_plan()
1189            .map_err(verification_failure)
1190    }
1191
1192    fn upgrade_governance(
1193        &self,
1194        plan: &DatastoreCreationPlan,
1195        binding: &DatastoreIdentityBinding,
1196        administrator: &ResolvedOperationSecret,
1197        catalog_credential: &ResolvedOperationSecret,
1198        storage: &[ResolvedCreationStorageSecret],
1199        batch: u32,
1200    ) -> Result<GovernanceUpgradeProgress, DatastoreCreationFailure> {
1201        ConfiguredDatastoreCreationProvisioner::upgrade_governance(
1202            self,
1203            plan,
1204            binding,
1205            administrator,
1206            catalog_credential,
1207            storage,
1208            batch,
1209        )
1210    }
1211
1212    fn seed_demo(
1213        &self,
1214        plan: &DatastoreCreationPlan,
1215        binding: &DatastoreIdentityBinding,
1216        administrator: &ResolvedOperationSecret,
1217        catalog_credential: &ResolvedOperationSecret,
1218        storage: &[ResolvedCreationStorageSecret],
1219    ) -> Result<(), DatastoreCreationFailure> {
1220        let (expected_path, _) = creation_lake_location(plan).map_err(verification_failure)?;
1221        let scratch = create_scratch(plan, std::iter::once(expected_path.clone()))?;
1222        let executor_root = plan
1223            .dataset_executor_root()
1224            .ok_or_else(|| verification_failure(()))?;
1225        let executor = ahri_tre_lake::DatasetExecutor::acquire(
1226            std::path::Path::new(executor_root),
1227            binding.datastore_id,
1228        )
1229        .map_err(verification_failure)?;
1230        let maintenance = creation_postgres_connection(plan)
1231            .dataset_maintenance(plan.postgresql().database(), administrator.material())
1232            .map_err(verification_failure)?;
1233        let mut exclusion = maintenance
1234            .authorize_namespace_reset(executor)
1235            .map_err(verification_failure)?;
1236        let opened = open_existing_lake(plan, catalog_credential.material(), storage, &scratch)?;
1237        if binding.datastore_name != plan.postgresql().database() {
1238            return Err(verification_failure(()));
1239        }
1240
1241        let study_id = StudyId(Uuid::parse_str(DEMO_STUDY_ID).map_err(verification_failure)?);
1242        let adapter = DuckLakeAdapter::new(expected_path);
1243        let seeds = reviewed_demo_dataset_seeds().map_err(verification_failure)?;
1244        adapter
1245            .replace_server_managed_dataset_tables_from_csv(
1246                &opened,
1247                &scratch,
1248                study_id,
1249                &seeds,
1250                &mut exclusion,
1251            )
1252            .map(|_| ())
1253            .map_err(lake_failure)
1254    }
1255
1256    fn apply_catalog_credential_rotation(
1257        &self,
1258        plan: &DatastoreCreationPlan,
1259        binding: &DatastoreIdentityBinding,
1260        administrator: &ResolvedOperationSecret,
1261        prior: &ResolvedOperationSecret,
1262        staged: &StagedCatalogCredential,
1263        rotated_at: DateTime<Utc>,
1264    ) -> Result<DatastoreIdentityBinding, DatastoreCreationFailure> {
1265        let version = i32::try_from(staged.version()).map_err(|_| postgresql_failure())?;
1266        apply_datastore_credential_rotation(
1267            &creation_postgres_connection(plan),
1268            administrator.material(),
1269            prior.material(),
1270            staged.material(),
1271            plan.postgresql().database(),
1272            plan.catalog_database(),
1273            plan.catalog_role(),
1274            binding,
1275            version,
1276            rotated_at,
1277        )
1278        .map_err(|error| match error {
1279            PgDatastoreCredentialRotationError::Restored => postgresql_failure(),
1280            PgDatastoreCredentialRotationError::StateUncertain => credential_state_uncertain(),
1281        })
1282    }
1283
1284    fn rollback_catalog_credential_rotation(
1285        &self,
1286        plan: &DatastoreCreationPlan,
1287        binding: &DatastoreIdentityBinding,
1288        administrator: &ResolvedOperationSecret,
1289        prior: &ResolvedOperationSecret,
1290    ) -> Result<(), DatastoreCreationFailure> {
1291        rollback_datastore_credential_rotation(
1292            &creation_postgres_connection(plan),
1293            administrator.material(),
1294            prior.material(),
1295            plan.postgresql().database(),
1296            plan.catalog_role(),
1297            plan.catalog_database(),
1298            binding,
1299        )
1300        .map_err(|error| match error {
1301            PgDatastoreCredentialRotationError::Restored => postgresql_failure(),
1302            PgDatastoreCredentialRotationError::StateUncertain => credential_state_uncertain(),
1303        })
1304    }
1305}
1306
1307fn creation_postgres_connection(plan: &DatastoreCreationPlan) -> PgDatastoreCreationConnection {
1308    let root_certificates = match plan.postgresql().tls() {
1309        SessionTls::System => Vec::new(),
1310        SessionTls::CustomCa { certificates } => certificates.clone(),
1311    };
1312    PgDatastoreCreationConnection::new(
1313        plan.postgresql().host(),
1314        plan.postgresql().port(),
1315        plan.postgresql().bootstrap_database(),
1316        plan.administrator().username(),
1317        root_certificates,
1318    )
1319}
1320
1321fn postgresql_failure() -> DatastoreCreationFailure {
1322    DatastoreCreationFailure::new(
1323        DatastoreCreationFailureCode::PostgresqlSetupFailed,
1324        "PostgreSQL Datastore initialization did not complete",
1325    )
1326}
1327
1328fn lake_failure(_error: impl std::fmt::Debug) -> DatastoreCreationFailure {
1329    DatastoreCreationFailure::new(
1330        DatastoreCreationFailureCode::LakeSetupFailed,
1331        "Lake initialization did not complete",
1332    )
1333}
1334
1335fn verification_failure(_error: impl std::fmt::Debug) -> DatastoreCreationFailure {
1336    DatastoreCreationFailure::new(
1337        DatastoreCreationFailureCode::VerificationFailed,
1338        "Cross-store verification did not complete",
1339    )
1340}
1341
1342fn internal_failure() -> DatastoreCreationFailure {
1343    DatastoreCreationFailure::new(
1344        DatastoreCreationFailureCode::InternalFailed,
1345        "Datastore creation state could not be persisted",
1346    )
1347}
1348
1349fn credential_state_uncertain() -> DatastoreCreationFailure {
1350    DatastoreCreationFailure::new(
1351        DatastoreCreationFailureCode::CredentialStateUncertain,
1352        "Credential owner state requires reconciliation",
1353    )
1354}
1355
1356fn creation_lake_location(
1357    plan: &DatastoreCreationPlan,
1358) -> Result<(String, DatastoreLakeLocation), ()> {
1359    match plan.lake().storage() {
1360        LakeStoragePlan::Filesystem { base_path } => {
1361            let config = FilesystemLakeSessionConfig::new(
1362                base_path,
1363                plan.datastore_prefix(),
1364                plan.postgresql().host(),
1365                plan.postgresql().port(),
1366                plan.catalog_database(),
1367                plan.catalog_schema(),
1368                plan.catalog_role(),
1369                catalog_tls(plan.postgresql().tls()),
1370                EncryptionMode::None,
1371            )
1372            .map_err(|_| ())?;
1373            Ok((
1374                config.canonical_data_path(),
1375                DatastoreLakeLocation::Filesystem {
1376                    base_path: base_path.clone(),
1377                    directory: plan.datastore_prefix().to_string(),
1378                },
1379            ))
1380        }
1381        storage => {
1382            let location = object_creation_location(plan).map_err(|_| ())?;
1383            let lake_path = location.canonical_data_path();
1384            let persisted = match (storage, location) {
1385                (
1386                    LakeStoragePlan::AwsS3 { .. },
1387                    ObjectStorageLocation::AwsS3 {
1388                        region,
1389                        bucket,
1390                        prefix,
1391                    },
1392                ) => DatastoreLakeLocation::AwsS3 {
1393                    region,
1394                    bucket,
1395                    prefix,
1396                },
1397                (
1398                    LakeStoragePlan::S3Compatible { .. },
1399                    ObjectStorageLocation::S3Compatible {
1400                        endpoint,
1401                        url_style,
1402                        bucket,
1403                        prefix,
1404                    },
1405                ) => DatastoreLakeLocation::S3Compatible {
1406                    endpoint,
1407                    url_style: match url_style {
1408                        S3UrlStyle::Path => "path",
1409                        S3UrlStyle::VirtualHosted => "virtual_hosted",
1410                    }
1411                    .to_string(),
1412                    bucket,
1413                    prefix,
1414                },
1415                (
1416                    LakeStoragePlan::AzureBlob { .. },
1417                    ObjectStorageLocation::AzureBlob {
1418                        account_endpoint,
1419                        container,
1420                        prefix,
1421                    },
1422                ) => DatastoreLakeLocation::AzureBlob {
1423                    account_endpoint,
1424                    container,
1425                    prefix,
1426                },
1427                _ => return Err(()),
1428            };
1429            Ok((lake_path, persisted))
1430        }
1431    }
1432}
1433
1434fn object_creation_location(
1435    plan: &DatastoreCreationPlan,
1436) -> Result<ObjectStorageLocation, ahri_tre_lake::LakeError> {
1437    let location = match plan.lake().storage() {
1438        LakeStoragePlan::AwsS3 {
1439            region,
1440            bucket,
1441            base_prefix,
1442            ..
1443        } => ObjectStorageLocation::aws_s3_for_datastore(
1444            region,
1445            bucket,
1446            base_prefix.as_deref(),
1447            plan.datastore_prefix(),
1448        )?,
1449        LakeStoragePlan::S3Compatible {
1450            endpoint,
1451            url_style,
1452            bucket,
1453            base_prefix,
1454            ..
1455        } => ObjectStorageLocation::s3_compatible_for_datastore(
1456            endpoint,
1457            match url_style {
1458                ConfiguredS3UrlStyle::Path => S3UrlStyle::Path,
1459                ConfiguredS3UrlStyle::VirtualHosted => S3UrlStyle::VirtualHosted,
1460            },
1461            bucket,
1462            base_prefix.as_deref(),
1463            plan.datastore_prefix(),
1464        )?,
1465        LakeStoragePlan::AzureBlob {
1466            account_endpoint,
1467            container,
1468            base_prefix,
1469            ..
1470        } => ObjectStorageLocation::azure_blob_for_datastore(
1471            account_endpoint,
1472            container,
1473            base_prefix.as_deref(),
1474            plan.datastore_prefix(),
1475        )?,
1476        LakeStoragePlan::Filesystem { .. } => {
1477            return Err(ahri_tre_lake::LakeError::InvalidObjectStorageIdentity);
1478        }
1479    };
1480    Ok(location)
1481}
1482
1483fn create_scratch(
1484    plan: &DatastoreCreationPlan,
1485    lake_paths: impl IntoIterator<Item = String>,
1486) -> Result<ScratchAttempt, DatastoreCreationFailure> {
1487    let paths = lake_paths
1488        .into_iter()
1489        .map(std::path::PathBuf::from)
1490        .collect::<Vec<_>>();
1491    let scratch = TrustedScratch::open(
1492        std::path::Path::new(plan.configuration_scratch_root()),
1493        paths.iter().map(std::path::PathBuf::as_path),
1494    )
1495    .map_err(lake_failure)?;
1496    let id = uuid::Uuid::new_v4().simple().to_string();
1497    scratch
1498        .create_attempt(ScratchAttemptId::new(&id).map_err(lake_failure)?)
1499        .map_err(lake_failure)
1500}
1501
1502fn initialize_object_lake(
1503    config: &ObjectLakeSessionConfig,
1504    catalog_credential: &SecretMaterial,
1505    plan: &DatastoreCreationPlan,
1506    storage: &[ResolvedCreationStorageSecret],
1507    scratch: &ScratchAttempt,
1508) -> Result<DuckLakeOpenedCatalog, DatastoreCreationFailure> {
1509    let resolved = |configured: &ConfiguredSecret| {
1510        storage
1511            .iter()
1512            .find(|secret| secret.configured() == configured)
1513            .map(ResolvedCreationStorageSecret::resolved)
1514            .ok_or_else(|| lake_failure(()))
1515    };
1516    let authentication = plan
1517        .lake()
1518        .storage()
1519        .authentication()
1520        .ok_or_else(|| lake_failure(()))?;
1521    match authentication {
1522        StorageAuthenticationPlan::S3Static {
1523            access_key_id,
1524            secret_access_key,
1525        } => {
1526            let access = resolved(access_key_id)?;
1527            let secret = resolved(secret_access_key)?;
1528            DuckLakeAdapter::initialize_object_storage_catalog(
1529                config,
1530                catalog_credential,
1531                ObjectStorageAuthentication::S3Static {
1532                    access_key_id: access.material(),
1533                    secret_access_key: secret.material(),
1534                },
1535                scratch,
1536            )
1537        }
1538        StorageAuthenticationPlan::AwsWorkloadIdentity { source } => {
1539            let source = match source {
1540                ConfiguredAwsWorkloadIdentitySource::EcsTask => AwsWorkloadIdentitySource::EcsTask,
1541                ConfiguredAwsWorkloadIdentitySource::Ec2Instance => {
1542                    AwsWorkloadIdentitySource::Ec2Instance
1543                }
1544            };
1545            DuckLakeAdapter::initialize_object_storage_catalog(
1546                config,
1547                catalog_credential,
1548                ObjectStorageAuthentication::AwsWorkloadIdentity { source },
1549                scratch,
1550            )
1551        }
1552        StorageAuthenticationPlan::AzureManagedIdentity { identity } => {
1553            let client_id = match identity {
1554                AzureManagedIdentity::SystemAssigned => None,
1555                AzureManagedIdentity::UserAssigned { client_id } => Some(*client_id),
1556            };
1557            DuckLakeAdapter::initialize_object_storage_catalog(
1558                config,
1559                catalog_credential,
1560                ObjectStorageAuthentication::AzureManagedIdentity { client_id },
1561                scratch,
1562            )
1563        }
1564        StorageAuthenticationPlan::AzureServicePrincipal {
1565            tenant_id,
1566            client_id,
1567            client_secret,
1568        } => {
1569            let secret = resolved(client_secret)?;
1570            DuckLakeAdapter::initialize_object_storage_catalog(
1571                config,
1572                catalog_credential,
1573                ObjectStorageAuthentication::AzureServicePrincipal {
1574                    tenant_id: *tenant_id,
1575                    client_id: *client_id,
1576                    client_secret: secret.material(),
1577                },
1578                scratch,
1579            )
1580        }
1581    }
1582    .map_err(lake_failure)
1583}
1584
1585fn open_existing_lake(
1586    plan: &DatastoreCreationPlan,
1587    catalog_credential: &SecretMaterial,
1588    storage: &[ResolvedCreationStorageSecret],
1589    scratch: &ScratchAttempt,
1590) -> Result<DuckLakeOpenedCatalog, DatastoreCreationFailure> {
1591    let catalog_tls = catalog_tls(plan.postgresql().tls());
1592    match plan.lake().storage() {
1593        LakeStoragePlan::Filesystem { base_path } => {
1594            let config = FilesystemLakeSessionConfig::new(
1595                base_path,
1596                plan.datastore_prefix(),
1597                plan.postgresql().host(),
1598                plan.postgresql().port(),
1599                plan.catalog_database(),
1600                plan.catalog_schema(),
1601                plan.catalog_role(),
1602                catalog_tls,
1603                EncryptionMode::ServerManaged,
1604            )
1605            .map_err(lake_failure)?;
1606            DuckLakeAdapter::open_existing_filesystem_catalog_with_scratch(
1607                &config,
1608                catalog_credential,
1609                scratch,
1610            )
1611            .map_err(lake_failure)
1612        }
1613        _ => {
1614            let location = object_creation_location(plan).map_err(lake_failure)?;
1615            let config = ObjectLakeSessionConfig::new(
1616                location,
1617                plan.postgresql().host(),
1618                plan.postgresql().port(),
1619                plan.catalog_database(),
1620                plan.catalog_schema(),
1621                plan.catalog_role(),
1622                catalog_tls,
1623                object_storage_tls(plan.lake().storage()),
1624                EncryptionMode::ServerManaged,
1625            )
1626            .map_err(lake_failure)?;
1627            open_existing_object_lake(&config, catalog_credential, plan, storage, scratch)
1628        }
1629    }
1630}
1631
1632fn open_existing_object_lake(
1633    config: &ObjectLakeSessionConfig,
1634    catalog_credential: &SecretMaterial,
1635    plan: &DatastoreCreationPlan,
1636    storage: &[ResolvedCreationStorageSecret],
1637    scratch: &ScratchAttempt,
1638) -> Result<DuckLakeOpenedCatalog, DatastoreCreationFailure> {
1639    let resolved = |configured: &ConfiguredSecret| {
1640        storage
1641            .iter()
1642            .find(|secret| secret.configured() == configured)
1643            .map(ResolvedCreationStorageSecret::resolved)
1644            .ok_or_else(|| lake_failure(()))
1645    };
1646    let authentication = plan
1647        .lake()
1648        .storage()
1649        .authentication()
1650        .ok_or_else(|| lake_failure(()))?;
1651    match authentication {
1652        StorageAuthenticationPlan::S3Static {
1653            access_key_id,
1654            secret_access_key,
1655        } => {
1656            let access = resolved(access_key_id)?;
1657            let secret = resolved(secret_access_key)?;
1658            DuckLakeAdapter::open_existing_object_storage_catalog(
1659                config,
1660                catalog_credential,
1661                ObjectStorageAuthentication::S3Static {
1662                    access_key_id: access.material(),
1663                    secret_access_key: secret.material(),
1664                },
1665                scratch,
1666            )
1667        }
1668        StorageAuthenticationPlan::AwsWorkloadIdentity { source } => {
1669            let source = match source {
1670                ConfiguredAwsWorkloadIdentitySource::EcsTask => AwsWorkloadIdentitySource::EcsTask,
1671                ConfiguredAwsWorkloadIdentitySource::Ec2Instance => {
1672                    AwsWorkloadIdentitySource::Ec2Instance
1673                }
1674            };
1675            DuckLakeAdapter::open_existing_object_storage_catalog(
1676                config,
1677                catalog_credential,
1678                ObjectStorageAuthentication::AwsWorkloadIdentity { source },
1679                scratch,
1680            )
1681        }
1682        StorageAuthenticationPlan::AzureManagedIdentity { identity } => {
1683            let client_id = match identity {
1684                AzureManagedIdentity::SystemAssigned => None,
1685                AzureManagedIdentity::UserAssigned { client_id } => Some(*client_id),
1686            };
1687            DuckLakeAdapter::open_existing_object_storage_catalog(
1688                config,
1689                catalog_credential,
1690                ObjectStorageAuthentication::AzureManagedIdentity { client_id },
1691                scratch,
1692            )
1693        }
1694        StorageAuthenticationPlan::AzureServicePrincipal {
1695            tenant_id,
1696            client_id,
1697            client_secret,
1698        } => {
1699            let secret = resolved(client_secret)?;
1700            DuckLakeAdapter::open_existing_object_storage_catalog(
1701                config,
1702                catalog_credential,
1703                ObjectStorageAuthentication::AzureServicePrincipal {
1704                    tenant_id: *tenant_id,
1705                    client_id: *client_id,
1706                    client_secret: secret.material(),
1707                },
1708                scratch,
1709            )
1710        }
1711    }
1712    .map_err(lake_failure)
1713}
1714
1715fn catalog_tls(tls: &SessionTls) -> FilesystemCatalogTls {
1716    match tls {
1717        SessionTls::System => FilesystemCatalogTls::System,
1718        SessionTls::CustomCa { certificates } => FilesystemCatalogTls::CustomCa {
1719            certificates: certificates.clone(),
1720        },
1721    }
1722}
1723
1724fn object_storage_tls(storage: &LakeStoragePlan) -> ObjectStorageTls {
1725    match storage {
1726        LakeStoragePlan::S3Compatible { tls, .. } => match tls {
1727            StorageTlsPlan::System => ObjectStorageTls::System,
1728            StorageTlsPlan::CustomCa { certificates } => ObjectStorageTls::CustomCa {
1729                certificates: certificates.clone(),
1730            },
1731        },
1732        _ => ObjectStorageTls::System,
1733    }
1734}
1735
1736enum RemovalReferenceAuthority<'a, A, P> {
1737    Application(&'a A),
1738    DatastoreBindings {
1739        authority: &'a A,
1740        provisioner: &'a P,
1741    },
1742    SessionRecords(&'a A),
1743    AuthenticationArtifacts(&'a A),
1744}
1745
1746impl<A, P> ManagedSecretReferenceAuthority for RemovalReferenceAuthority<'_, A, P>
1747where
1748    A: DatastoreCreationAuthority,
1749    P: DatastoreCreationProvisioner,
1750{
1751    fn reference_status(
1752        &self,
1753        reference: &ManagedSecretReference,
1754    ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
1755        match self {
1756            Self::Application(authority) => authority.application_reference_status(reference),
1757            Self::SessionRecords(authority) => authority.session_reference_status(reference),
1758            Self::AuthenticationArtifacts(authority) => {
1759                authority.authentication_reference_status(reference)
1760            }
1761            Self::DatastoreBindings {
1762                authority,
1763                provisioner,
1764            } => {
1765                let durable_status = authority.datastore_reference_status(reference)?;
1766                let targets = authority
1767                    .postgresql_binding_inspection_targets()
1768                    .map_err(|_| ManagedSecretReferenceInspectionError::unavailable())?;
1769                let mut bindings = Vec::new();
1770                for target in targets {
1771                    let administrator = authority
1772                        .resolve_postgresql_binding_inspection_administrator(&target)
1773                        .map_err(|_| ManagedSecretReferenceInspectionError::unavailable())?;
1774                    bindings.extend(
1775                        provisioner
1776                            .inspect_all_bindings(&target, &administrator)
1777                            .map_err(|_| ManagedSecretReferenceInspectionError::unavailable())?,
1778                    );
1779                }
1780                if bindings.iter().any(|binding| {
1781                    binding.managed_secret_ref.as_deref()
1782                        == Some(reference.as_secret_reference().as_str())
1783                }) {
1784                    return Ok(ManagedSecretReferenceStatus::Referenced);
1785                }
1786                if durable_status == ManagedSecretReferenceStatus::Referenced {
1787                    for datastore_id in
1788                        authority.reconcilable_datastore_binding_reference_ids(reference)?
1789                    {
1790                        if !bindings
1791                            .iter()
1792                            .any(|binding| binding.datastore_id == datastore_id)
1793                        {
1794                            authority.release_absent_datastore_binding_reference(
1795                                datastore_id,
1796                                reference,
1797                            )?;
1798                        }
1799                    }
1800                }
1801                authority.datastore_reference_status(reference)
1802            }
1803        }
1804    }
1805}
1806
1807/// Authenticated handler that turns logical create intent into one durable saga.
1808pub struct PrivilegedDatastoreCreationHandler<A, P> {
1809    authority: A,
1810    provisioner: P,
1811}
1812
1813impl<A, P> PrivilegedDatastoreCreationHandler<A, P> {
1814    pub fn new(authority: A, provisioner: P) -> Self {
1815        Self {
1816            authority,
1817            provisioner,
1818        }
1819    }
1820}
1821
1822impl<A, P> AdministrationRequestHandler for PrivilegedDatastoreCreationHandler<A, P>
1823where
1824    A: DatastoreCreationAuthority,
1825    P: DatastoreCreationProvisioner,
1826{
1827    fn handle(
1828        &self,
1829        peer: AdministrationPeerIdentity,
1830        request: AdministrationRequest,
1831    ) -> AdministrationResponse {
1832        match request {
1833            AdministrationRequest::DatastoreCreate { datastore } => self.create(peer, &datastore),
1834            AdministrationRequest::DatastoreReconcile { datastore } => {
1835                self.reconcile(peer, &datastore)
1836            }
1837            AdministrationRequest::DatastoreUpgradeGovernance {
1838                datastore,
1839                batch,
1840                plan_only,
1841            } => self.upgrade_governance(peer, &datastore, batch, plan_only),
1842            AdministrationRequest::DatastoreSeedDemo { datastore } => {
1843                self.seed_demo(peer, &datastore)
1844            }
1845            AdministrationRequest::DatastoreRotateCredential { datastore } => {
1846                self.rotate_credential(peer, &datastore)
1847            }
1848            AdministrationRequest::ManagedSecretRemove {
1849                deployment_id,
1850                reference,
1851            } => self.remove_secret(peer, deployment_id, &reference),
1852        }
1853    }
1854}
1855
1856impl<A, P> PrivilegedDatastoreCreationHandler<A, P>
1857where
1858    A: DatastoreCreationAuthority,
1859    P: DatastoreCreationProvisioner,
1860{
1861    fn remove_secret(
1862        &self,
1863        peer: AdministrationPeerIdentity,
1864        deployment_id: Uuid,
1865        reference: &ManagedSecretReference,
1866    ) -> AdministrationResponse {
1867        let failed = |code: AdministrationFailureCode| {
1868            audit_secret_outcome(peer, reference, code.as_str());
1869            AdministrationResponse::managed_secret_removal_failed(reference.as_str(), code)
1870        };
1871        if self.authority.deployment_id() != Ok(deployment_id) {
1872            return failed(AdministrationFailureCode::ConfigurationUnavailable);
1873        }
1874        let application: RemovalReferenceAuthority<'_, A, P> =
1875            RemovalReferenceAuthority::Application(&self.authority);
1876        let datastores = RemovalReferenceAuthority::DatastoreBindings {
1877            authority: &self.authority,
1878            provisioner: &self.provisioner,
1879        };
1880        let sessions: RemovalReferenceAuthority<'_, A, P> =
1881            RemovalReferenceAuthority::SessionRecords(&self.authority);
1882        let authentication: RemovalReferenceAuthority<'_, A, P> =
1883            RemovalReferenceAuthority::AuthenticationArtifacts(&self.authority);
1884        let proof = match ManagedSecretReferenceAuthorities::new(
1885            &application,
1886            &datastores,
1887            &sessions,
1888            &authentication,
1889        )
1890        .prove_unreferenced(reference)
1891        {
1892            Ok(proof) => proof,
1893            Err(error) => {
1894                return failed(match error.kind() {
1895                    ManagedSecretRemovalProofErrorKind::Referenced => {
1896                        AdministrationFailureCode::ReferencePresent
1897                    }
1898                    ManagedSecretRemovalProofErrorKind::InspectionUnavailable => {
1899                        AdministrationFailureCode::ReferenceInspectionUnavailable
1900                    }
1901                });
1902            }
1903        };
1904        let removed = match self.authority.remove_managed_secret(reference, proof) {
1905            Ok(metadata) => metadata,
1906            Err(_) => return failed(AdministrationFailureCode::SecretRemovalFailed),
1907        };
1908        audit_secret_outcome(peer, reference, "removed");
1909        AdministrationResponse::managed_secret_removed(reference.as_str(), removed.version().get())
1910    }
1911    fn upgrade_governance(
1912        &self,
1913        peer: AdministrationPeerIdentity,
1914        datastore: &DatastoreConfigurationId,
1915        batch: u32,
1916        plan_only: bool,
1917    ) -> AdministrationResponse {
1918        let failed = |code: AdministrationFailureCode| {
1919            audit_outcome(
1920                peer,
1921                "datastore_upgrade_governance",
1922                datastore,
1923                code.as_str(),
1924            );
1925            AdministrationResponse::datastore_failed(datastore.as_str(), code)
1926        };
1927        if batch == 0 || batch > 1000 {
1928            return failed(AdministrationFailureCode::ConfigurationUnavailable);
1929        }
1930        let run = || -> Result<_, AdministrationFailureCode> {
1931            let initial = self.authority.plan(datastore)?;
1932            let administrator = self.authority.resolve_administrator(&initial)?;
1933            let binding = self
1934                .provisioner
1935                .find_existing(&initial, &administrator)
1936                .map_err(|e| e.administration_code())?
1937                .ok_or(AdministrationFailureCode::VerificationFailed)?;
1938            let plan = self
1939                .authority
1940                .plan_for_existing(datastore, binding.datastore_id)?;
1941            let catalog = self
1942                .authority
1943                .resolve_existing_catalog_credential(&plan, &binding)?;
1944            let storage = self.authority.resolve_storage(&plan)?;
1945            if plan_only {
1946                if plan.governance_custodians().is_empty() {
1947                    return Err(AdministrationFailureCode::ConfigurationUnavailable);
1948                }
1949                let p = self
1950                    .provisioner
1951                    .plan_governance(&plan, &binding, &administrator, &catalog, &storage)
1952                    .map_err(|e| e.administration_code())?;
1953                return Ok(AdministrationResponse::GovernancePlanned {
1954                    datastore: datastore.as_str().into(),
1955                    datastore_id: binding.datastore_id,
1956                    prepared: p.prepared,
1957                    unclassified_assets: p.unclassified_assets,
1958                    immutable_versions: p.immutable_versions,
1959                    custodians: plan.governance_custodians().to_vec(),
1960                });
1961            }
1962            let progress = self
1963                .provisioner
1964                .upgrade_governance(&plan, &binding, &administrator, &catalog, &storage, batch)
1965                .map_err(|e| e.administration_code())?;
1966            Ok(AdministrationResponse::GovernanceUpgraded {
1967                datastore: datastore.as_str().into(),
1968                datastore_id: binding.datastore_id,
1969                converted_assets: progress.converted_assets,
1970                remaining_unclassified_assets: progress.remaining_unclassified_assets,
1971                bound_datasets: progress.bound_datasets,
1972                bound_datafiles: progress.bound_datafiles,
1973                remaining_dataset_bindings: progress.remaining_dataset_bindings,
1974                remaining_datafile_bindings: progress.remaining_datafile_bindings,
1975                remaining_evidence: progress.remaining_evidence,
1976            })
1977        };
1978        match run() {
1979            Ok(response) => {
1980                audit_outcome(
1981                    peer,
1982                    "datastore_upgrade_governance",
1983                    datastore,
1984                    if plan_only {
1985                        "planned"
1986                    } else {
1987                        "batch_complete"
1988                    },
1989                );
1990                response
1991            }
1992            Err(code) => failed(code),
1993        }
1994    }
1995
1996    fn seed_demo(
1997        &self,
1998        peer: AdministrationPeerIdentity,
1999        datastore: &DatastoreConfigurationId,
2000    ) -> AdministrationResponse {
2001        let failed = |code: AdministrationFailureCode| {
2002            audit_outcome(peer, "datastore_seed_demo", datastore, code.as_str());
2003            AdministrationResponse::datastore_failed(datastore.as_str(), code)
2004        };
2005        if datastore.as_str() != DEMO_DATASTORE {
2006            return failed(AdministrationFailureCode::ConfigurationUnavailable);
2007        }
2008        if let response @ AdministrationResponse::Failed { .. } = self.reconcile(peer, datastore) {
2009            return response;
2010        }
2011        let initial_plan = match self.authority.plan(datastore) {
2012            Ok(plan) => plan,
2013            Err(code) => return failed(code),
2014        };
2015        let administrator = match self.authority.resolve_administrator(&initial_plan) {
2016            Ok(secret) => secret,
2017            Err(code) => return failed(code),
2018        };
2019        let binding = match self
2020            .provisioner
2021            .find_existing(&initial_plan, &administrator)
2022        {
2023            Ok(Some(binding)) => binding,
2024            Ok(None) => return failed(AdministrationFailureCode::VerificationFailed),
2025            Err(error) => return failed(error.administration_code()),
2026        };
2027        let plan = match self
2028            .authority
2029            .plan_for_existing(datastore, binding.datastore_id)
2030        {
2031            Ok(plan) => plan,
2032            Err(code) => return failed(code),
2033        };
2034        let catalog_credential = match self
2035            .authority
2036            .resolve_existing_catalog_credential(&plan, &binding)
2037        {
2038            Ok(secret) => secret,
2039            Err(code) => return failed(code),
2040        };
2041        let storage = match self.authority.resolve_storage(&plan) {
2042            Ok(storage) => storage,
2043            Err(code) => return failed(code),
2044        };
2045        if let Err(error) = self.provisioner.seed_demo(
2046            &plan,
2047            &binding,
2048            &administrator,
2049            &catalog_credential,
2050            &storage,
2051        ) {
2052            return failed(error.administration_code());
2053        }
2054        audit_outcome(peer, "datastore_seed_demo", datastore, "ready");
2055        AdministrationResponse::datastore_ready(datastore.as_str(), binding.datastore_id)
2056    }
2057
2058    fn reconcile(
2059        &self,
2060        peer: AdministrationPeerIdentity,
2061        datastore: &DatastoreConfigurationId,
2062    ) -> AdministrationResponse {
2063        let failed = |code: AdministrationFailureCode| {
2064            audit_outcome(peer, "datastore_reconcile", datastore, code.as_str());
2065            AdministrationResponse::datastore_failed(datastore.as_str(), code)
2066        };
2067        let initial_plan = match self.authority.plan(datastore) {
2068            Ok(plan) => plan,
2069            Err(code) => return failed(code),
2070        };
2071        let administrator = match self.authority.resolve_administrator(&initial_plan) {
2072            Ok(secret) => secret,
2073            Err(code) => return failed(code),
2074        };
2075        let existing = match self
2076            .provisioner
2077            .find_existing(&initial_plan, &administrator)
2078        {
2079            Ok(existing) => existing,
2080            Err(error) => return failed(error.administration_code()),
2081        };
2082        let Some(binding) = existing else {
2083            return self.create(peer, datastore);
2084        };
2085        let plan = match self
2086            .authority
2087            .plan_for_existing(datastore, binding.datastore_id)
2088        {
2089            Ok(plan) => plan,
2090            Err(code) => return failed(code),
2091        };
2092        let catalog_credential = match self
2093            .authority
2094            .resolve_existing_catalog_credential(&plan, &binding)
2095        {
2096            Ok(secret) => secret,
2097            Err(code) => return failed(code),
2098        };
2099        let storage = match self.authority.resolve_storage(&plan) {
2100            Ok(storage) => storage,
2101            Err(code) => return failed(code),
2102        };
2103        if let Err(error) = self.provisioner.verify_existing(
2104            &plan,
2105            &binding,
2106            &administrator,
2107            &catalog_credential,
2108            &storage,
2109        ) {
2110            return failed(error.administration_code());
2111        }
2112        if let Err(code) = self
2113            .authority
2114            .retain_datastore_binding_reference(&plan, plan.catalog_secret_reference())
2115        {
2116            return failed(code);
2117        }
2118        audit_outcome(peer, "datastore_reconcile", datastore, "ready");
2119        AdministrationResponse::datastore_ready(datastore.as_str(), binding.datastore_id)
2120    }
2121
2122    fn rotate_credential(
2123        &self,
2124        peer: AdministrationPeerIdentity,
2125        datastore: &DatastoreConfigurationId,
2126    ) -> AdministrationResponse {
2127        let failed = |code: AdministrationFailureCode| {
2128            audit_outcome(
2129                peer,
2130                "datastore_rotate_credential",
2131                datastore,
2132                code.as_str(),
2133            );
2134            AdministrationResponse::datastore_failed(datastore.as_str(), code)
2135        };
2136        let initial_plan = match self.authority.plan(datastore) {
2137            Ok(plan) => plan,
2138            Err(code) => return failed(code),
2139        };
2140        let administrator = match self.authority.resolve_administrator(&initial_plan) {
2141            Ok(secret) => secret,
2142            Err(code) => return failed(code),
2143        };
2144        let binding = match self
2145            .provisioner
2146            .find_existing(&initial_plan, &administrator)
2147        {
2148            Ok(Some(binding)) if binding.lifecycle_state == DatastoreLifecycleState::Ready => {
2149                binding
2150            }
2151            Ok(_) => return failed(AdministrationFailureCode::VerificationFailed),
2152            Err(error) => return failed(error.administration_code()),
2153        };
2154        let plan = match self
2155            .authority
2156            .plan_for_existing(datastore, binding.datastore_id)
2157        {
2158            Ok(plan) => plan,
2159            Err(code) => return failed(code),
2160        };
2161        if let Err(code) = self
2162            .authority
2163            .retain_datastore_binding_reference(&plan, plan.catalog_secret_reference())
2164        {
2165            return failed(code);
2166        }
2167        let prior = match self
2168            .authority
2169            .resolve_existing_catalog_credential(&plan, &binding)
2170        {
2171            Ok(secret) => secret,
2172            Err(code) => return failed(code),
2173        };
2174        let generated = match generate_catalog_credential() {
2175            Ok(secret) => secret,
2176            Err(()) => return failed(AdministrationFailureCode::InternalFailed),
2177        };
2178        let staged = match self
2179            .authority
2180            .stage_catalog_credential_replacement(&plan, &binding, &generated)
2181        {
2182            Ok(staged) => staged,
2183            Err(code) => return failed(code),
2184        };
2185        let rotated_at = Utc::now();
2186        let externally_applied = match self.provisioner.apply_catalog_credential_rotation(
2187            &plan,
2188            &binding,
2189            &administrator,
2190            &prior,
2191            &staged,
2192            rotated_at,
2193        ) {
2194            Ok(binding) => binding,
2195            Err(error) => {
2196                if error.administration_code()
2197                    == AdministrationFailureCode::CredentialStateUncertain
2198                {
2199                    return failed(AdministrationFailureCode::CredentialStateUncertain);
2200                }
2201                if self
2202                    .authority
2203                    .discard_catalog_credential_replacement(&plan, &staged)
2204                    .is_err()
2205                {
2206                    return failed(AdministrationFailureCode::InternalFailed);
2207                }
2208                return failed(error.administration_code());
2209            }
2210        };
2211        if i32::try_from(staged.version()).ok() != Some(externally_applied.credential_version) {
2212            let rollback = self.provisioner.rollback_catalog_credential_rotation(
2213                &plan,
2214                &binding,
2215                &administrator,
2216                &prior,
2217            );
2218            if rollback.is_err() {
2219                return failed(AdministrationFailureCode::CredentialStateUncertain);
2220            }
2221            return failed(
2222                if self
2223                    .authority
2224                    .discard_catalog_credential_replacement(&plan, &staged)
2225                    .is_ok()
2226                {
2227                    AdministrationFailureCode::VerificationFailed
2228                } else {
2229                    AdministrationFailureCode::InternalFailed
2230                },
2231            );
2232        }
2233        if let Err(code) = self
2234            .authority
2235            .activate_catalog_credential_replacement(&plan, &staged)
2236        {
2237            if code == AdministrationFailureCode::CredentialStateUncertain {
2238                return failed(code);
2239            }
2240            let rollback = self.provisioner.rollback_catalog_credential_rotation(
2241                &plan,
2242                &binding,
2243                &administrator,
2244                &prior,
2245            );
2246            let discard = self
2247                .authority
2248                .discard_catalog_credential_replacement(&plan, &staged);
2249            return failed(if rollback.is_ok() && discard.is_ok() {
2250                code
2251            } else {
2252                AdministrationFailureCode::CredentialStateUncertain
2253            });
2254        }
2255        audit_outcome(peer, "datastore_rotate_credential", datastore, "ready");
2256        AdministrationResponse::datastore_ready(datastore.as_str(), binding.datastore_id)
2257    }
2258
2259    fn create(
2260        &self,
2261        peer: AdministrationPeerIdentity,
2262        datastore: &DatastoreConfigurationId,
2263    ) -> AdministrationResponse {
2264        let failed = |code: AdministrationFailureCode| {
2265            audit_outcome(peer, "datastore_create", datastore, code.as_str());
2266            AdministrationResponse::datastore_failed(datastore.as_str(), code)
2267        };
2268        let plan = match self.authority.plan(datastore) {
2269            Ok(plan) => plan,
2270            Err(code) => return failed(code),
2271        };
2272        let administrator = match self.authority.resolve_administrator(&plan) {
2273            Ok(secret) => secret,
2274            Err(code) => return failed(code),
2275        };
2276        let storage = match self.authority.resolve_storage(&plan) {
2277            Ok(storage) => storage,
2278            Err(code) => return failed(code),
2279        };
2280        let catalog_credential = match generate_catalog_credential() {
2281            Ok(secret) => secret,
2282            Err(()) => return failed(AdministrationFailureCode::InternalFailed),
2283        };
2284
2285        let mut pending = match self
2286            .provisioner
2287            .begin(&plan, &administrator, &catalog_credential)
2288        {
2289            DatastoreCreationStart::Started(pending) => pending,
2290            DatastoreCreationStart::Conflict => {
2291                return failed(AdministrationFailureCode::Conflict);
2292            }
2293            DatastoreCreationStart::Failed(error) => {
2294                return failed(error.administration_code());
2295            }
2296            DatastoreCreationStart::FailedWithBinding {
2297                mut pending,
2298                failure,
2299            } => {
2300                if self
2301                    .provisioner
2302                    .mark_failed(&mut pending, &failure)
2303                    .is_err()
2304                {
2305                    return failed(AdministrationFailureCode::InternalFailed);
2306                }
2307                return failed(failure.administration_code());
2308            }
2309        };
2310        let created = match self
2311            .authority
2312            .create_catalog_credential(&plan, &catalog_credential)
2313        {
2314            Ok(metadata) => metadata,
2315            Err(_) => {
2316                let error = DatastoreCreationFailure::new(
2317                    DatastoreCreationFailureCode::CredentialSetupFailed,
2318                    "Managed catalog credential was not activated",
2319                );
2320                if self.provisioner.mark_failed(&mut pending, &error).is_err() {
2321                    return failed(AdministrationFailureCode::InternalFailed);
2322                }
2323                return failed(error.administration_code());
2324            }
2325        };
2326
2327        if self
2328            .authority
2329            .verify_catalog_credential(&plan, &catalog_credential, &created)
2330            .is_err()
2331        {
2332            let error = DatastoreCreationFailure::new(
2333                DatastoreCreationFailureCode::CredentialSetupFailed,
2334                "Managed catalog credential verification did not complete",
2335            );
2336            if self.provisioner.mark_failed(&mut pending, &error).is_err() {
2337                return failed(AdministrationFailureCode::InternalFailed);
2338            }
2339            return failed(error.administration_code());
2340        }
2341
2342        if let Err(error) =
2343            self.provisioner
2344                .initialize_lake(&mut pending, &plan, &catalog_credential, &storage)
2345        {
2346            if self.provisioner.mark_failed(&mut pending, &error).is_err() {
2347                return failed(AdministrationFailureCode::InternalFailed);
2348            }
2349            return failed(error.administration_code());
2350        }
2351        if let Err(error) = self.provisioner.verify(&mut pending, &plan) {
2352            if self.provisioner.mark_failed(&mut pending, &error).is_err() {
2353                return failed(AdministrationFailureCode::InternalFailed);
2354            }
2355            return failed(error.administration_code());
2356        }
2357        if let Err(code) = self
2358            .authority
2359            .retain_datastore_binding_reference(&plan, plan.catalog_secret_reference())
2360        {
2361            let error = DatastoreCreationFailure::new(
2362                DatastoreCreationFailureCode::CredentialSetupFailed,
2363                "durable Datastore credential ownership could not be recorded",
2364            );
2365            if self.provisioner.mark_failed(&mut pending, &error).is_err() {
2366                return failed(AdministrationFailureCode::InternalFailed);
2367            }
2368            return failed(code);
2369        }
2370        if let Err(error) = self.provisioner.mark_ready(&mut pending, &created) {
2371            if self.provisioner.mark_failed(&mut pending, &error).is_err() {
2372                return failed(AdministrationFailureCode::InternalFailed);
2373            }
2374            return failed(error.administration_code());
2375        }
2376
2377        audit_outcome(peer, "datastore_create", datastore, "ready");
2378        AdministrationResponse::datastore_ready(datastore.as_str(), plan.datastore_id())
2379    }
2380}
2381
2382fn reviewed_demo_dataset_seeds() -> Result<Vec<DatasetCsvSeed<'static>>, ()> {
2383    const FIXTURES: [(&str, &[u8]); 12] = [
2384        ("entities", include_bytes!("../../../data/entities.csv")),
2385        ("episodes", include_bytes!("../../../data/episodes.csv")),
2386        (
2387            "household_map",
2388            include_bytes!("../../../data/household_map.csv"),
2389        ),
2390        (
2391            "householdmemberships",
2392            include_bytes!("../../../data/householdmemberships.csv"),
2393        ),
2394        (
2395            "householdresidences",
2396            include_bytes!("../../../data/householdresidences.csv"),
2397        ),
2398        ("households", include_bytes!("../../../data/households.csv")),
2399        (
2400            "individual_map",
2401            include_bytes!("../../../data/individual_map.csv"),
2402        ),
2403        (
2404            "individualresidences",
2405            include_bytes!("../../../data/individualresidences.csv"),
2406        ),
2407        (
2408            "individuals",
2409            include_bytes!("../../../data/individuals.csv"),
2410        ),
2411        (
2412            "location_map",
2413            include_bytes!("../../../data/location_map.csv"),
2414        ),
2415        ("locations", include_bytes!("../../../data/locations.csv")),
2416        (
2417            "transformations",
2418            include_bytes!("../../../data/transformations.csv"),
2419        ),
2420    ];
2421
2422    FIXTURES
2423        .into_iter()
2424        .map(|(name, contents)| {
2425            Ok(DatasetCsvSeed {
2426                dataset_name: NcName::parse(name).map_err(|_| ())?,
2427                contents,
2428            })
2429        })
2430        .collect()
2431}
2432
2433fn audit_secret_outcome(
2434    peer: AdministrationPeerIdentity,
2435    reference: &ManagedSecretReference,
2436    outcome: &'static str,
2437) {
2438    tracing::info!(
2439        peer_uid = peer.uid(),
2440        peer_gid = peer.gid(),
2441        peer_pid = peer.pid(),
2442        command = "managed_secret_remove",
2443        reference = reference.as_str(),
2444        outcome,
2445        "privileged administration command completed"
2446    );
2447}
2448fn audit_outcome(
2449    peer: AdministrationPeerIdentity,
2450    command: &'static str,
2451    datastore: &DatastoreConfigurationId,
2452    outcome: &'static str,
2453) {
2454    tracing::info!(
2455        peer_uid = peer.uid(),
2456        peer_gid = peer.gid(),
2457        peer_pid = peer.pid(),
2458        command,
2459        datastore = datastore.as_str(),
2460        outcome,
2461        "privileged administration command completed"
2462    );
2463}
2464
2465fn generate_catalog_credential() -> Result<SecretMaterial, ()> {
2466    let encoded = format!(
2467        "tre-{}-{}-{}",
2468        uuid::Uuid::new_v4().simple(),
2469        uuid::Uuid::new_v4().simple(),
2470        uuid::Uuid::new_v4().simple()
2471    );
2472    SecretMaterial::try_from(encoded.into_bytes()).map_err(|_| ())
2473}
2474
2475fn equivalent_ducklake_data_paths(left: &str, right: &str) -> bool {
2476    left == right || left.trim_end_matches('/') == right.trim_end_matches('/')
2477}
2478
2479/// One resumable operator batch. Safe metadata remains available while any of
2480/// these prerequisites is pending; ordinary Session opens cannot invoke it.
2481#[derive(Debug)]
2482pub struct GovernanceUpgradeProgress {
2483    pub converted_assets: u64,
2484    pub remaining_unclassified_assets: u64,
2485    pub bound_datasets: usize,
2486    pub bound_datafiles: usize,
2487    pub remaining_datafile_bindings: bool,
2488    pub remaining_dataset_bindings: bool,
2489    pub remaining_evidence: bool,
2490}
2491impl ConfiguredDatastoreCreationProvisioner {
2492    pub fn upgrade_governance(
2493        &self,
2494        plan: &DatastoreCreationPlan,
2495        binding: &DatastoreIdentityBinding,
2496        administrator: &ResolvedOperationSecret,
2497        catalog_credential: &ResolvedOperationSecret,
2498        storage: &[ResolvedCreationStorageSecret],
2499        batch: u32,
2500    ) -> Result<GovernanceUpgradeProgress, DatastoreCreationFailure> {
2501        if batch == 0 || batch > 1000 {
2502            return Err(verification_failure(()));
2503        }
2504        self.verify_existing(plan, binding, administrator, catalog_credential, storage)?;
2505        let coordinator = plan
2506            .dataset_executor_root()
2507            .ok_or_else(|| verification_failure(()))?;
2508        let executor = ahri_tre_lake::DatasetExecutor::acquire(
2509            std::path::Path::new(coordinator),
2510            binding.datastore_id,
2511        )
2512        .map_err(verification_failure)?;
2513        let maintenance = creation_postgres_connection(plan)
2514            .dataset_maintenance(plan.postgresql().database(), administrator.material())
2515            .map_err(verification_failure)?;
2516        let _exclusion = maintenance
2517            .authorize_namespace_reset(executor)
2518            .map_err(verification_failure)?;
2519
2520        let appointment = ahri_tre_pgmeta::governance::GovernanceAppointment::from_identities(
2521            plan.governance_custodians()
2522                .iter()
2523                .map(|a| (a.issuer(), a.subject())),
2524        )
2525        .map_err(verification_failure)?;
2526        let operator = creation_postgres_connection(plan)
2527            .upgrade_governance(
2528                plan.postgresql().database(),
2529                plan.metadata_role(),
2530                administrator.material(),
2531                &appointment,
2532            )
2533            .map_err(verification_failure)?;
2534        let (path, _) = creation_lake_location(plan).map_err(verification_failure)?;
2535        let scratch = create_scratch(plan, std::iter::once(path))?;
2536        let mut lake = open_existing_lake(plan, catalog_credential.material(), storage, &scratch)?;
2537        let ledger = operator.binding().map_err(verification_failure)?;
2538        ahri_tre_lake::GovernanceLedger::provision(&mut lake.connection, StudyId(ledger.study_id))
2539            .map_err(verification_failure)?;
2540        let converted_assets = operator
2541            .convert_existing(batch)
2542            .map_err(verification_failure)?;
2543        let origins = operator
2544            .unbound_legacy_datasets(batch)
2545            .map_err(verification_failure)?;
2546        let bound_datasets = origins.len();
2547        for origin in origins {
2548            let (table_uuid, lake_snapshot) = ahri_tre_lake::dataset_physical_identity(
2549                &lake.connection,
2550                plan.catalog_schema(),
2551                StudyId(origin.study_id),
2552                origin.name.as_str(),
2553                origin.major,
2554                origin.minor,
2555                origin.patch,
2556            )
2557            .map_err(verification_failure)?;
2558            operator
2559                .bind_dataset(&ahri_tre_types::GovernedDatasetLocation {
2560                    version_id: origin.version_id,
2561                    asset_id: origin.asset_id,
2562                    study_id: origin.study_id,
2563                    name: origin.name,
2564                    major: origin.major,
2565                    minor: origin.minor,
2566                    patch: origin.patch,
2567                    table_uuid,
2568                    lake_snapshot,
2569                })
2570                .map_err(verification_failure)?;
2571        }
2572        let files = operator
2573            .unbound_legacy_datafiles(batch)
2574            .map_err(verification_failure)?;
2575        let bound_datafiles = files.len();
2576        for origin in files {
2577            let attestation = ahri_tre_lake::attest_datafile(
2578                &lake.attach_plan.data_path,
2579                StudyId(origin.study_id),
2580                ahri_tre_types::AssetId(origin.asset_id),
2581                &origin.datafile,
2582            )
2583            .map_err(verification_failure)?;
2584            operator
2585                .bind_datafile(&attestation)
2586                .map_err(verification_failure)?;
2587        }
2588        for event in operator.pending(batch).map_err(verification_failure)? {
2589            crate::disclosure::project_session_evidence(&operator, &mut lake.connection, event)
2590                .map_err(verification_failure)?;
2591        }
2592        Ok(GovernanceUpgradeProgress {
2593            converted_assets,
2594            remaining_unclassified_assets: operator
2595                .upgrade_plan()
2596                .map_err(verification_failure)?
2597                .unclassified_assets,
2598            bound_datasets,
2599            bound_datafiles,
2600            remaining_datafile_bindings: !operator
2601                .unbound_legacy_datafiles(1)
2602                .map_err(verification_failure)?
2603                .is_empty(),
2604            remaining_dataset_bindings: !operator
2605                .unbound_legacy_datasets(1)
2606                .map_err(verification_failure)?
2607                .is_empty(),
2608            remaining_evidence: !operator
2609                .pending(1)
2610                .map_err(verification_failure)?
2611                .is_empty(),
2612        })
2613    }
2614}
2615
2616#[cfg(test)]
2617mod tests {
2618    use super::{equivalent_ducklake_data_paths, reviewed_demo_dataset_seeds};
2619
2620    #[test]
2621    fn ducklake_persisted_trailing_separator_is_equivalent() {
2622        assert!(equivalent_ducklake_data_paths(
2623            "/var/lib/ahri-tre/lake/development/",
2624            "/var/lib/ahri-tre/lake/development"
2625        ));
2626        assert!(!equivalent_ducklake_data_paths(
2627            "/var/lib/ahri-tre/lake/development-other",
2628            "/var/lib/ahri-tre/lake/development"
2629        ));
2630    }
2631
2632    #[test]
2633    fn trusted_demo_fixtures_are_the_exact_reviewed_allowlist() {
2634        let fixtures = reviewed_demo_dataset_seeds().unwrap();
2635        assert_eq!(
2636            fixtures
2637                .iter()
2638                .map(|fixture| fixture.dataset_name.as_str())
2639                .collect::<Vec<_>>(),
2640            vec![
2641                "entities",
2642                "episodes",
2643                "household_map",
2644                "householdmemberships",
2645                "householdresidences",
2646                "households",
2647                "individual_map",
2648                "individualresidences",
2649                "individuals",
2650                "location_map",
2651                "locations",
2652                "transformations",
2653            ]
2654        );
2655        assert!(fixtures.iter().all(|fixture| !fixture.contents.is_empty()));
2656    }
2657}
2658#[cfg(test)]
2659mod ticket19_tests {
2660    use super::{
2661        CreatedCatalogCredential, DatastoreCreationAuthority, DatastoreCreationFailure,
2662        DatastoreCreationProvisioner, DatastoreCreationStart, PrivilegedDatastoreCreationHandler,
2663        ResolvedCreationStorageSecret, StagedCatalogCredential,
2664    };
2665    use age::secrecy::ExposeSecret;
2666    use ahri_tre_config::{
2667        ConfigurationAuthority, DatastoreConfigurationId, EffectiveConfigurationTarget,
2668    };
2669    use ahri_tre_runtime::{
2670        AdministrationFailureCode, AdministrationPeerIdentity, AdministrationRequest,
2671        AdministrationRequestHandler, AdministrationResponse, DatastoreCreationPlan,
2672        ResolvedOperationSecret, ResolvedSecretVersion,
2673    };
2674    use ahri_tre_secrets::internal::{
2675        ManagedSecretOwnerBackend, ManagedSecretReferenceInspectionError,
2676        ManagedSecretReferenceStatus, ManagedSecretRemovalProof, ManagedSecretTrustedBackend,
2677    };
2678    use ahri_tre_secrets::{
2679        LocalManagedSecretStore, ManagedSecretBackend, ManagedSecretErrorKind,
2680        ManagedSecretIdentity, ManagedSecretMetadata, ManagedSecretReference, SecretMaterial,
2681    };
2682    use ahri_tre_types::{
2683        DatastoreCreationFailureCode, DatastoreIdentityBinding, DatastoreLifecycleState,
2684        LakeCatalogCredentialMode,
2685    };
2686    use chrono::Utc;
2687    use std::fs;
2688    use std::os::unix::fs::PermissionsExt;
2689    use std::path::PathBuf;
2690    use std::str::FromStr;
2691    use std::sync::{Arc, Mutex};
2692    use uuid::Uuid;
2693
2694    struct RotationAuthority {
2695        plan: DatastoreCreationPlan,
2696        store: Arc<LocalManagedSecretStore>,
2697        events: Arc<Mutex<Vec<&'static str>>>,
2698        application_referenced: bool,
2699    }
2700
2701    impl DatastoreCreationAuthority for RotationAuthority {
2702        fn deployment_id(&self) -> Result<Uuid, AdministrationFailureCode> {
2703            Ok(Uuid::from_u128(1))
2704        }
2705
2706        fn retain_datastore_binding_reference(
2707            &self,
2708            _plan: &DatastoreCreationPlan,
2709            _reference: &ManagedSecretReference,
2710        ) -> Result<(), AdministrationFailureCode> {
2711            Ok(())
2712        }
2713
2714        fn plan(
2715            &self,
2716            _datastore: &DatastoreConfigurationId,
2717        ) -> Result<DatastoreCreationPlan, AdministrationFailureCode> {
2718            Ok(self.plan.clone())
2719        }
2720
2721        fn resolve_administrator(
2722            &self,
2723            _plan: &DatastoreCreationPlan,
2724        ) -> Result<ResolvedOperationSecret, AdministrationFailureCode> {
2725            Ok(resolved_secret(
2726                "managed://postgresql/runtime-password",
2727                b"administrator",
2728                1,
2729            ))
2730        }
2731
2732        fn plan_for_existing(
2733            &self,
2734            _datastore: &DatastoreConfigurationId,
2735            _datastore_id: Uuid,
2736        ) -> Result<DatastoreCreationPlan, AdministrationFailureCode> {
2737            Ok(self.plan.clone())
2738        }
2739
2740        fn resolve_existing_catalog_credential(
2741            &self,
2742            plan: &DatastoreCreationPlan,
2743            binding: &DatastoreIdentityBinding,
2744        ) -> Result<ResolvedOperationSecret, AdministrationFailureCode> {
2745            self.events.lock().unwrap().push("resolve-prior");
2746            Ok(resolved_secret(
2747                plan.catalog_secret_reference().as_str(),
2748                b"prior-catalog-credential",
2749                u64::try_from(binding.credential_version).unwrap(),
2750            ))
2751        }
2752
2753        fn resolve_storage(
2754            &self,
2755            _plan: &DatastoreCreationPlan,
2756        ) -> Result<Vec<ResolvedCreationStorageSecret>, AdministrationFailureCode> {
2757            Ok(Vec::new())
2758        }
2759
2760        fn create_catalog_credential(
2761            &self,
2762            _plan: &DatastoreCreationPlan,
2763            _material: &SecretMaterial,
2764        ) -> Result<CreatedCatalogCredential, AdministrationFailureCode> {
2765            unreachable!("rotation never creates a second logical reference")
2766        }
2767
2768        fn verify_catalog_credential(
2769            &self,
2770            _plan: &DatastoreCreationPlan,
2771            _material: &SecretMaterial,
2772            _credential: &CreatedCatalogCredential,
2773        ) -> Result<(), AdministrationFailureCode> {
2774            unreachable!("rotation verifies through the owning PostgreSQL adapter")
2775        }
2776
2777        fn stage_catalog_credential_replacement(
2778            &self,
2779            plan: &DatastoreCreationPlan,
2780            _binding: &DatastoreIdentityBinding,
2781            material: &SecretMaterial,
2782        ) -> Result<StagedCatalogCredential, AdministrationFailureCode> {
2783            self.events.lock().unwrap().push("stage");
2784            let active = self
2785                .store
2786                .lookup_active_metadata(plan.catalog_secret_reference())
2787                .map_err(|_| AdministrationFailureCode::CredentialSetupFailed)?;
2788            self.store
2789                .stage_owner_replacement(
2790                    plan.catalog_secret_reference(),
2791                    material,
2792                    active.version(),
2793                )
2794                .map(|inner| StagedCatalogCredential { inner })
2795                .map_err(|_| AdministrationFailureCode::CredentialSetupFailed)
2796        }
2797
2798        fn activate_catalog_credential_replacement(
2799            &self,
2800            _plan: &DatastoreCreationPlan,
2801            staged: &StagedCatalogCredential,
2802        ) -> Result<CreatedCatalogCredential, AdministrationFailureCode> {
2803            self.events.lock().unwrap().push("activate");
2804            self.store
2805                .activate_owner_replacement(&staged.inner)
2806                .map(Into::into)
2807                .map_err(|_| AdministrationFailureCode::CredentialSetupFailed)
2808        }
2809
2810        fn discard_catalog_credential_replacement(
2811            &self,
2812            _plan: &DatastoreCreationPlan,
2813            staged: &StagedCatalogCredential,
2814        ) -> Result<(), AdministrationFailureCode> {
2815            self.events.lock().unwrap().push("discard");
2816            self.store
2817                .discard_owner_replacement(&staged.inner)
2818                .map_err(|_| AdministrationFailureCode::CredentialSetupFailed)
2819        }
2820
2821        fn configured_datastore_ids(
2822            &self,
2823        ) -> Result<Vec<DatastoreConfigurationId>, AdministrationFailureCode> {
2824            Ok(Vec::new())
2825        }
2826
2827        fn postgresql_binding_inspection_targets(
2828            &self,
2829        ) -> Result<
2830            Vec<ahri_tre_runtime::PostgresqlBindingInspectionTarget>,
2831            AdministrationFailureCode,
2832        > {
2833            Ok(Vec::new())
2834        }
2835
2836        fn application_reference_status(
2837            &self,
2838            _reference: &ManagedSecretReference,
2839        ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
2840            self.events.lock().unwrap().push("application-inspect");
2841            Ok(if self.application_referenced {
2842                ManagedSecretReferenceStatus::Referenced
2843            } else {
2844                ManagedSecretReferenceStatus::Unreferenced
2845            })
2846        }
2847
2848        fn session_reference_status(
2849            &self,
2850            _reference: &ManagedSecretReference,
2851        ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
2852            Ok(ManagedSecretReferenceStatus::Unreferenced)
2853        }
2854
2855        fn datastore_reference_status(
2856            &self,
2857            reference: &ManagedSecretReference,
2858        ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
2859            self.store
2860                .trusted_held_reference_status(reference)
2861                .map_err(|_| ManagedSecretReferenceInspectionError::unavailable())
2862        }
2863
2864        fn authentication_reference_status(
2865            &self,
2866            _reference: &ManagedSecretReference,
2867        ) -> Result<ManagedSecretReferenceStatus, ManagedSecretReferenceInspectionError> {
2868            Ok(ManagedSecretReferenceStatus::Unreferenced)
2869        }
2870
2871        fn remove_managed_secret(
2872            &self,
2873            reference: &ManagedSecretReference,
2874            proof: ManagedSecretRemovalProof,
2875        ) -> Result<ManagedSecretMetadata, AdministrationFailureCode> {
2876            self.store
2877                .remove_with_repository_evidence(reference, proof)
2878                .map_err(|_| AdministrationFailureCode::CredentialSetupFailed)
2879        }
2880    }
2881
2882    struct RotationProvisioner {
2883        binding: DatastoreIdentityBinding,
2884        events: Arc<Mutex<Vec<&'static str>>>,
2885        fail_apply: bool,
2886        uncertain_apply: bool,
2887    }
2888
2889    impl DatastoreCreationProvisioner for RotationProvisioner {
2890        type Pending = ();
2891
2892        fn begin(
2893            &self,
2894            _plan: &DatastoreCreationPlan,
2895            _administrator: &ResolvedOperationSecret,
2896            _catalog_credential: &SecretMaterial,
2897        ) -> DatastoreCreationStart<Self::Pending> {
2898            unreachable!("rotation does not run Datastore creation")
2899        }
2900
2901        fn initialize_lake(
2902            &self,
2903            _pending: &mut Self::Pending,
2904            _plan: &DatastoreCreationPlan,
2905            _catalog_credential: &SecretMaterial,
2906            _storage: &[ResolvedCreationStorageSecret],
2907        ) -> Result<(), DatastoreCreationFailure> {
2908            unreachable!("rotation does not initialize the lake")
2909        }
2910
2911        fn verify(
2912            &self,
2913            _pending: &mut Self::Pending,
2914            _plan: &DatastoreCreationPlan,
2915        ) -> Result<(), DatastoreCreationFailure> {
2916            unreachable!("rotation verifies in its owner-specific phase")
2917        }
2918
2919        fn mark_ready(
2920            &self,
2921            _pending: &mut Self::Pending,
2922            _credential: &CreatedCatalogCredential,
2923        ) -> Result<(), DatastoreCreationFailure> {
2924            unreachable!("rotation updates the retained ready binding")
2925        }
2926
2927        fn mark_failed(
2928            &self,
2929            _pending: &mut Self::Pending,
2930            _failure: &DatastoreCreationFailure,
2931        ) -> Result<(), DatastoreCreationFailure> {
2932            unreachable!("rotation compensates instead of failing the Datastore")
2933        }
2934
2935        fn find_existing(
2936            &self,
2937            _plan: &DatastoreCreationPlan,
2938            _administrator: &ResolvedOperationSecret,
2939        ) -> Result<Option<DatastoreIdentityBinding>, DatastoreCreationFailure> {
2940            Ok(Some(self.binding.clone()))
2941        }
2942
2943        fn apply_catalog_credential_rotation(
2944            &self,
2945            _plan: &DatastoreCreationPlan,
2946            _binding: &DatastoreIdentityBinding,
2947            _administrator: &ResolvedOperationSecret,
2948            _prior: &ResolvedOperationSecret,
2949            staged: &StagedCatalogCredential,
2950            rotated_at: chrono::DateTime<Utc>,
2951        ) -> Result<DatastoreIdentityBinding, DatastoreCreationFailure> {
2952            self.events.lock().unwrap().push("apply-and-verify");
2953            if self.uncertain_apply {
2954                return Err(DatastoreCreationFailure::new(
2955                    DatastoreCreationFailureCode::CredentialStateUncertain,
2956                    "injected uncertain owner state",
2957                ));
2958            }
2959            if self.fail_apply {
2960                return Err(DatastoreCreationFailure::new(
2961                    DatastoreCreationFailureCode::PostgresqlSetupFailed,
2962                    "injected owner apply failure",
2963                ));
2964            }
2965            let mut binding = self.binding.clone();
2966            binding.credential_version = i32::try_from(staged.version()).unwrap();
2967            binding.credential_last_rotated_at = Some(rotated_at);
2968            binding.binding_fingerprint = binding.computed_binding_fingerprint();
2969            Ok(binding)
2970        }
2971
2972        fn rollback_catalog_credential_rotation(
2973            &self,
2974            _plan: &DatastoreCreationPlan,
2975            _binding: &DatastoreIdentityBinding,
2976            _administrator: &ResolvedOperationSecret,
2977            _prior: &ResolvedOperationSecret,
2978        ) -> Result<(), DatastoreCreationFailure> {
2979            self.events.lock().unwrap().push("rollback");
2980            Ok(())
2981        }
2982    }
2983
2984    #[test]
2985    fn owner_rotation_activates_only_after_external_apply_and_verification() {
2986        let fixture = ManagedStoreFixture::new("activate");
2987        let store = Arc::new(fixture.initialize_store());
2988        let plan = rotation_plan();
2989        let reference = plan.catalog_secret_reference().clone();
2990        store
2991            .create(
2992                &reference,
2993                &SecretMaterial::try_from(b"prior-catalog-credential".to_vec()).unwrap(),
2994            )
2995            .unwrap();
2996        let events = Arc::new(Mutex::new(Vec::new()));
2997        let binding = ready_binding(plan.datastore_id(), &reference);
2998        let handler = PrivilegedDatastoreCreationHandler::new(
2999            RotationAuthority {
3000                plan,
3001                store: Arc::clone(&store),
3002                events: Arc::clone(&events),
3003                application_referenced: false,
3004            },
3005            RotationProvisioner {
3006                binding: binding.clone(),
3007                events: Arc::clone(&events),
3008                fail_apply: false,
3009                uncertain_apply: false,
3010            },
3011        );
3012
3013        let response = handler.handle(peer(), rotate_request());
3014
3015        assert_eq!(
3016            response,
3017            AdministrationResponse::datastore_ready("research", binding.datastore_id)
3018        );
3019        assert_eq!(
3020            events.lock().unwrap().as_slice(),
3021            ["resolve-prior", "stage", "apply-and-verify", "activate"]
3022        );
3023        assert_eq!(
3024            store
3025                .resolve(&reference)
3026                .unwrap()
3027                .metadata()
3028                .version()
3029                .get(),
3030            2
3031        );
3032        assert!(
3033            store
3034                .resume_owner_replacement(&reference)
3035                .unwrap()
3036                .is_none()
3037        );
3038    }
3039
3040    #[test]
3041    fn failed_owner_apply_discards_stage_and_preserves_prior_active_version() {
3042        let fixture = ManagedStoreFixture::new("compensate");
3043        let store = Arc::new(fixture.initialize_store());
3044        let plan = rotation_plan();
3045        let reference = plan.catalog_secret_reference().clone();
3046        store
3047            .create(
3048                &reference,
3049                &SecretMaterial::try_from(b"prior-catalog-credential".to_vec()).unwrap(),
3050            )
3051            .unwrap();
3052        let events = Arc::new(Mutex::new(Vec::new()));
3053        let binding = ready_binding(plan.datastore_id(), &reference);
3054        let handler = PrivilegedDatastoreCreationHandler::new(
3055            RotationAuthority {
3056                plan,
3057                store: Arc::clone(&store),
3058                events: Arc::clone(&events),
3059                application_referenced: false,
3060            },
3061            RotationProvisioner {
3062                binding,
3063                events: Arc::clone(&events),
3064                fail_apply: true,
3065                uncertain_apply: false,
3066            },
3067        );
3068
3069        let response = handler.handle(peer(), rotate_request());
3070
3071        assert_eq!(
3072            response,
3073            AdministrationResponse::datastore_failed(
3074                "research",
3075                AdministrationFailureCode::PostgresqlSetupFailed,
3076            )
3077        );
3078        assert_eq!(
3079            events.lock().unwrap().as_slice(),
3080            ["resolve-prior", "stage", "apply-and-verify", "discard"]
3081        );
3082        assert_eq!(
3083            store
3084                .resolve(&reference)
3085                .unwrap()
3086                .metadata()
3087                .version()
3088                .get(),
3089            1
3090        );
3091        assert!(
3092            store
3093                .resume_owner_replacement(&reference)
3094                .unwrap()
3095                .is_none()
3096        );
3097    }
3098    #[test]
3099    fn uncertain_owner_apply_retains_stage_for_authoritative_reconciliation() {
3100        let fixture = ManagedStoreFixture::new("uncertain");
3101        let store = Arc::new(fixture.initialize_store());
3102        let plan = rotation_plan();
3103        let reference = plan.catalog_secret_reference().clone();
3104        store
3105            .create(
3106                &reference,
3107                &SecretMaterial::try_from(b"prior-catalog-credential".to_vec()).unwrap(),
3108            )
3109            .unwrap();
3110        let events = Arc::new(Mutex::new(Vec::new()));
3111        let binding = ready_binding(plan.datastore_id(), &reference);
3112        let handler = PrivilegedDatastoreCreationHandler::new(
3113            RotationAuthority {
3114                plan,
3115                store: Arc::clone(&store),
3116                events: Arc::clone(&events),
3117                application_referenced: false,
3118            },
3119            RotationProvisioner {
3120                binding,
3121                events: Arc::clone(&events),
3122                fail_apply: false,
3123                uncertain_apply: true,
3124            },
3125        );
3126
3127        let response = handler.handle(peer(), rotate_request());
3128
3129        assert_eq!(
3130            response,
3131            AdministrationResponse::datastore_failed(
3132                "research",
3133                AdministrationFailureCode::CredentialStateUncertain,
3134            )
3135        );
3136        assert_eq!(
3137            events.lock().unwrap().as_slice(),
3138            ["resolve-prior", "stage", "apply-and-verify"]
3139        );
3140        assert_eq!(
3141            store
3142                .resolve(&reference)
3143                .unwrap()
3144                .metadata()
3145                .version()
3146                .get(),
3147            1
3148        );
3149        assert!(
3150            store
3151                .resume_owner_replacement(&reference)
3152                .unwrap()
3153                .is_some()
3154        );
3155    }
3156
3157    #[test]
3158    fn durable_datastore_marker_blocks_removal_after_configuration_entry_is_gone() {
3159        let fixture = ManagedStoreFixture::new("stale-datastore-binding");
3160        let store = Arc::new(fixture.initialize_store());
3161        let plan = rotation_plan();
3162        let reference =
3163            ManagedSecretReference::from_str("managed://datastore/orphan/credential").unwrap();
3164        store
3165            .create(
3166                &reference,
3167                &SecretMaterial::try_from(b"orphan-binding-credential".to_vec()).unwrap(),
3168            )
3169            .unwrap();
3170        store
3171            .replace_trusted_reference_holder(
3172                &format!("datastore:{}", Uuid::new_v4()),
3173                std::slice::from_ref(&reference),
3174            )
3175            .unwrap();
3176        let binding = ready_binding(plan.datastore_id(), plan.catalog_secret_reference());
3177        let handler = PrivilegedDatastoreCreationHandler::new(
3178            RotationAuthority {
3179                plan,
3180                store: Arc::clone(&store),
3181                events: Arc::new(Mutex::new(Vec::new())),
3182                application_referenced: false,
3183            },
3184            RotationProvisioner {
3185                binding,
3186                events: Arc::new(Mutex::new(Vec::new())),
3187                fail_apply: false,
3188                uncertain_apply: false,
3189            },
3190        );
3191
3192        let response = handler.handle(
3193            peer(),
3194            AdministrationRequest::ManagedSecretRemove {
3195                reference: reference.clone(),
3196                deployment_id: Uuid::from_u128(1),
3197            },
3198        );
3199
3200        assert_eq!(
3201            response,
3202            AdministrationResponse::managed_secret_removal_failed(
3203                reference.as_str(),
3204                AdministrationFailureCode::ReferencePresent,
3205            )
3206        );
3207        assert_eq!(
3208            store
3209                .resolve(&reference)
3210                .unwrap()
3211                .metadata()
3212                .version()
3213                .get(),
3214            1
3215        );
3216    }
3217
3218    #[test]
3219    fn removal_rejects_deployment_mismatch_before_reference_inspection() {
3220        let fixture = ManagedStoreFixture::new("deployment-mismatch");
3221        let store = Arc::new(fixture.initialize_store());
3222        let plan = rotation_plan();
3223        let reference =
3224            ManagedSecretReference::from_str("managed://operator/removal-candidate").unwrap();
3225        store
3226            .create(
3227                &reference,
3228                &SecretMaterial::try_from(b"removal-candidate".to_vec()).unwrap(),
3229            )
3230            .unwrap();
3231        let events = Arc::new(Mutex::new(Vec::new()));
3232        let binding = ready_binding(plan.datastore_id(), plan.catalog_secret_reference());
3233        let handler = PrivilegedDatastoreCreationHandler::new(
3234            RotationAuthority {
3235                plan,
3236                store: Arc::clone(&store),
3237                events: Arc::clone(&events),
3238                application_referenced: false,
3239            },
3240            RotationProvisioner {
3241                binding,
3242                events: Arc::clone(&events),
3243                fail_apply: false,
3244                uncertain_apply: false,
3245            },
3246        );
3247
3248        let response = handler.handle(
3249            peer(),
3250            AdministrationRequest::ManagedSecretRemove {
3251                reference: reference.clone(),
3252                deployment_id: Uuid::from_u128(2),
3253            },
3254        );
3255
3256        assert_eq!(
3257            response,
3258            AdministrationResponse::managed_secret_removal_failed(
3259                reference.as_str(),
3260                AdministrationFailureCode::ConfigurationUnavailable,
3261            )
3262        );
3263        assert!(events.lock().unwrap().is_empty());
3264        assert_eq!(
3265            store
3266                .resolve(&reference)
3267                .unwrap()
3268                .metadata()
3269                .version()
3270                .get(),
3271            1
3272        );
3273    }
3274
3275    #[test]
3276    fn removal_requires_all_authorities_to_report_unreferenced() {
3277        let fixture = ManagedStoreFixture::new("removal-proof");
3278        let store = Arc::new(fixture.initialize_store());
3279        let plan = rotation_plan();
3280        let reference =
3281            ManagedSecretReference::from_str("managed://operator/removal-candidate").unwrap();
3282        store
3283            .create(
3284                &reference,
3285                &SecretMaterial::try_from(b"removal-candidate".to_vec()).unwrap(),
3286            )
3287            .unwrap();
3288        let binding = ready_binding(plan.datastore_id(), plan.catalog_secret_reference());
3289        let handler = PrivilegedDatastoreCreationHandler::new(
3290            RotationAuthority {
3291                plan,
3292                store: Arc::clone(&store),
3293                events: Arc::new(Mutex::new(Vec::new())),
3294                application_referenced: true,
3295            },
3296            RotationProvisioner {
3297                binding,
3298                events: Arc::new(Mutex::new(Vec::new())),
3299                fail_apply: false,
3300                uncertain_apply: false,
3301            },
3302        );
3303
3304        let response = handler.handle(
3305            peer(),
3306            AdministrationRequest::ManagedSecretRemove {
3307                reference: reference.clone(),
3308                deployment_id: Uuid::from_u128(1),
3309            },
3310        );
3311
3312        assert_eq!(
3313            response,
3314            AdministrationResponse::managed_secret_removal_failed(
3315                reference.as_str(),
3316                AdministrationFailureCode::ReferencePresent,
3317            )
3318        );
3319        assert_eq!(
3320            store
3321                .resolve(&reference)
3322                .unwrap()
3323                .metadata()
3324                .version()
3325                .get(),
3326            1
3327        );
3328    }
3329
3330    #[test]
3331    fn unreferenced_removal_is_durable_and_tombstones_the_reference() {
3332        let fixture = ManagedStoreFixture::new("remove");
3333        let store = Arc::new(fixture.initialize_store());
3334        let plan = rotation_plan();
3335        let reference =
3336            ManagedSecretReference::from_str("managed://operator/removal-candidate").unwrap();
3337        store
3338            .create(
3339                &reference,
3340                &SecretMaterial::try_from(b"removal-candidate".to_vec()).unwrap(),
3341            )
3342            .unwrap();
3343        let binding = ready_binding(plan.datastore_id(), plan.catalog_secret_reference());
3344        let handler = PrivilegedDatastoreCreationHandler::new(
3345            RotationAuthority {
3346                plan,
3347                store: Arc::clone(&store),
3348                events: Arc::new(Mutex::new(Vec::new())),
3349                application_referenced: false,
3350            },
3351            RotationProvisioner {
3352                binding,
3353                events: Arc::new(Mutex::new(Vec::new())),
3354                fail_apply: false,
3355                uncertain_apply: false,
3356            },
3357        );
3358
3359        let response = handler.handle(
3360            peer(),
3361            AdministrationRequest::ManagedSecretRemove {
3362                reference: reference.clone(),
3363                deployment_id: Uuid::from_u128(1),
3364            },
3365        );
3366
3367        assert_eq!(
3368            response,
3369            AdministrationResponse::managed_secret_removed(reference.as_str(), 1)
3370        );
3371        assert_eq!(
3372            store.resolve(&reference).unwrap_err().kind(),
3373            ManagedSecretErrorKind::Missing
3374        );
3375        assert_eq!(
3376            store
3377                .create(
3378                    &reference,
3379                    &SecretMaterial::try_from(b"cannot-recreate".to_vec()).unwrap(),
3380                )
3381                .unwrap_err()
3382                .kind(),
3383            ManagedSecretErrorKind::VersionConflict
3384        );
3385    }
3386
3387    fn rotation_plan() -> DatastoreCreationPlan {
3388        let document = fs::read_to_string(concat!(
3389            env!("CARGO_MANIFEST_DIR"),
3390            "/../../examples/config/application-v1.toml"
3391        ))
3392        .unwrap();
3393        let datastore = DatastoreConfigurationId::from_str("research").unwrap();
3394        let effective = ConfigurationAuthority::parse_application(&document)
3395            .unwrap()
3396            .resolve_effective(EffectiveConfigurationTarget::DatastoreCreate { datastore })
3397            .unwrap();
3398        DatastoreCreationPlan::try_from((
3399            &effective,
3400            Uuid::parse_str("12345678-9abc-4def-8123-456789abcdef").unwrap(),
3401        ))
3402        .unwrap()
3403    }
3404
3405    fn ready_binding(
3406        datastore_id: Uuid,
3407        reference: &ManagedSecretReference,
3408    ) -> DatastoreIdentityBinding {
3409        let mut binding = DatastoreIdentityBinding {
3410            datastore_id,
3411            datastore_name: "Research".to_string(),
3412            ducklake_catalog_database: "ducklake_catalog".to_string(),
3413            ducklake_catalog_schema: Some("catalog".to_string()),
3414            lake_path: "/srv/ahri-tre/lake/datasets".to_string(),
3415            storage_policy_id: Some("filesystem".to_string()),
3416            lake_location: None,
3417            ducklake_encryption: true,
3418            lake_catalog_credential_mode: LakeCatalogCredentialMode::ManagedLocal,
3419            lake_catalog_role_name: Some("catalog_owner".to_string()),
3420            managed_secret_ref: Some(reference.as_str().to_string()),
3421            credential_version: 1,
3422            credential_last_rotated_at: Some(Utc::now()),
3423            created_at: Utc::now(),
3424            creating_tool_version: Some(env!("CARGO_PKG_VERSION").to_string()),
3425            binding_fingerprint: String::new(),
3426            lifecycle_state: DatastoreLifecycleState::Ready,
3427            failure_code: None,
3428            failure_summary: None,
3429        };
3430        binding.binding_fingerprint = binding.computed_binding_fingerprint();
3431        binding
3432    }
3433
3434    fn resolved_secret(reference: &str, material: &[u8], version: u64) -> ResolvedOperationSecret {
3435        ResolvedOperationSecret::new(
3436            SecretMaterial::try_from(material.to_vec()).unwrap(),
3437            ResolvedSecretVersion::managed(
3438                ManagedSecretReference::from_str(reference).unwrap(),
3439                version,
3440            ),
3441        )
3442    }
3443
3444    fn peer() -> AdministrationPeerIdentity {
3445        AdministrationPeerIdentity::for_test(1000, 1000, Some(42))
3446    }
3447
3448    fn rotate_request() -> AdministrationRequest {
3449        AdministrationRequest::DatastoreRotateCredential {
3450            datastore: DatastoreConfigurationId::from_str("research").unwrap(),
3451        }
3452    }
3453
3454    struct ManagedStoreFixture {
3455        root: PathBuf,
3456        store_root: PathBuf,
3457        deployment_id: Uuid,
3458        identity: String,
3459    }
3460
3461    impl ManagedStoreFixture {
3462        fn new(label: &str) -> Self {
3463            let root = std::env::temp_dir().join(format!(
3464                "ahri-tre-app-ticket19-{label}-{}-{}",
3465                std::process::id(),
3466                Uuid::new_v4()
3467            ));
3468            fs::create_dir(&root).unwrap();
3469            fs::set_permissions(&root, fs::Permissions::from_mode(0o700)).unwrap();
3470            let identity = age::x25519::Identity::generate()
3471                .to_string()
3472                .expose_secret()
3473                .to_string();
3474            Self {
3475                store_root: root.join("secrets"),
3476                root,
3477                deployment_id: Uuid::new_v4(),
3478                identity,
3479            }
3480        }
3481
3482        fn initialize_store(&self) -> LocalManagedSecretStore {
3483            LocalManagedSecretStore::initialize_at_for_test(
3484                &self.store_root,
3485                self.deployment_id,
3486                self.parsed_identity(),
3487            )
3488            .unwrap()
3489        }
3490
3491        fn parsed_identity(&self) -> ManagedSecretIdentity {
3492            ManagedSecretIdentity::from_material(
3493                &SecretMaterial::try_from(self.identity.as_bytes().to_vec()).unwrap(),
3494            )
3495            .unwrap()
3496        }
3497    }
3498
3499    impl Drop for ManagedStoreFixture {
3500        fn drop(&mut self) {
3501            let _ = fs::remove_dir_all(&self.root);
3502        }
3503    }
3504}