1use 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#[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}
73pub 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
97pub 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
130pub 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
652pub trait DatastoreCreationProvisioner: Send + Sync {
654 type Pending;
655
656 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#[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
1807pub 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#[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}