1use ahri_tre_secrets::SecretMaterial;
4use ahri_tre_types::{DatastoreIdentityBinding, NewDatastoreIdentityBinding};
5
6use crate::{
7 DirectPostgresConnectConfig, PgMetadataAdapter, PgMetadataConnection,
8 datastore_identity_table_exists, insert_datastore_identity_binding,
9 read_datastore_identity_bindings,
10};
11
12#[derive(Debug, Clone)]
14pub struct PgDatastoreCreationConnection {
15 host: String,
16 port: u16,
17 bootstrap_database: String,
18 administrator_username: String,
19 root_certificates: Vec<String>,
20}
21
22impl PgDatastoreCreationConnection {
23 pub fn new(
24 host: impl Into<String>,
25 port: u16,
26 bootstrap_database: impl Into<String>,
27 administrator_username: impl Into<String>,
28 root_certificates: Vec<String>,
29 ) -> Self {
30 Self {
31 host: host.into(),
32 port,
33 bootstrap_database: bootstrap_database.into(),
34 administrator_username: administrator_username.into(),
35 root_certificates,
36 }
37 }
38
39 pub fn dataset_maintenance(
41 &self,
42 database: &str,
43 credential: &SecretMaterial,
44 ) -> Result<crate::PgDatasetMaintenance, crate::PgMetaError> {
45 credential
46 .expose(|bytes| {
47 let password =
48 std::str::from_utf8(bytes).map_err(|_| crate::PgMetaError::Decode {
49 field: "maintenance_credential",
50 value: String::new(),
51 message: "maintenance credential is unavailable".into(),
52 })?;
53 self.adapter(database, password).connect()
54 })
55 .and_then(crate::PgDatasetMaintenance::new)
56 }
57
58 fn adapter(&self, database: &str, password: &str) -> PgMetadataAdapter {
59 self.adapter_for(database, &self.administrator_username, password)
60 }
61
62 fn adapter_for(&self, database: &str, username: &str, password: &str) -> PgMetadataAdapter {
63 let config = DirectPostgresConnectConfig {
64 host: self.host.clone(),
65 port: self.port,
66 dbname: database.to_string(),
67 username: username.to_string(),
68 password: Some(password.to_string()),
69 sslmode: Some("verify-full".to_string()),
70 connect_timeout_secs: Some(10),
71 };
72 if self.root_certificates.is_empty() {
73 PgMetadataAdapter::direct(config)
74 } else {
75 PgMetadataAdapter::direct_with_ca(config, self.root_certificates.clone())
76 }
77 }
78}
79
80#[derive(Debug, Clone)]
82pub struct PgDatastoreCreationSpec {
83 target_database: String,
84 target_owner: String,
85 catalog_database: String,
86 catalog_schema: String,
87 catalog_role: String,
88 binding: NewDatastoreIdentityBinding,
89 governance: crate::governance::GovernanceAppointment,
90}
91
92impl PgDatastoreCreationSpec {
93 #[allow(clippy::too_many_arguments)]
94 pub fn new(
95 target_database: impl Into<String>,
96 target_owner: impl Into<String>,
97 catalog_database: impl Into<String>,
98 catalog_schema: impl Into<String>,
99 catalog_role: impl Into<String>,
100 binding: NewDatastoreIdentityBinding,
101 governance: crate::governance::GovernanceAppointment,
102 ) -> Self {
103 Self {
104 target_database: target_database.into(),
105 target_owner: target_owner.into(),
106 catalog_database: catalog_database.into(),
107 catalog_schema: catalog_schema.into(),
108 catalog_role: catalog_role.into(),
109 binding,
110 governance,
111 }
112 }
113
114 fn identifiers(&self) -> [&str; 5] {
115 [
116 &self.target_database,
117 &self.target_owner,
118 &self.catalog_database,
119 &self.catalog_schema,
120 &self.catalog_role,
121 ]
122 }
123}
124
125#[derive(Debug, Clone)]
127pub struct PgDatastoreValidationSpec {
128 target_database: String,
129 target_owner: String,
130 catalog_database: String,
131 catalog_schema: String,
132 catalog_role: String,
133}
134
135impl PgDatastoreValidationSpec {
136 pub fn new(
137 target_database: impl Into<String>,
138 target_owner: impl Into<String>,
139 catalog_database: impl Into<String>,
140 catalog_schema: impl Into<String>,
141 catalog_role: impl Into<String>,
142 ) -> Self {
143 Self {
144 target_database: target_database.into(),
145 target_owner: target_owner.into(),
146 catalog_database: catalog_database.into(),
147 catalog_schema: catalog_schema.into(),
148 catalog_role: catalog_role.into(),
149 }
150 }
151
152 fn identifiers(&self) -> [&str; 5] {
153 [
154 &self.target_database,
155 &self.target_owner,
156 &self.catalog_database,
157 &self.catalog_schema,
158 &self.catalog_role,
159 ]
160 }
161}
162
163pub struct PgPendingDatastoreCreation {
165 _administrator: PgMetadataConnection,
166 store: PgMetadataConnection,
167}
168
169impl std::fmt::Debug for PgPendingDatastoreCreation {
170 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
171 formatter
172 .debug_struct("PgPendingDatastoreCreation")
173 .field("administrator_lock_held", &true)
174 .field("store", &self.store)
175 .finish_non_exhaustive()
176 }
177}
178
179impl PgPendingDatastoreCreation {
180 pub fn store_mut(&mut self) -> &mut PgMetadataConnection {
181 &mut self.store
182 }
183}
184
185#[derive(Debug)]
187pub enum PgDatastoreCreationStart {
188 Started(PgPendingDatastoreCreation),
189 Conflict,
190 FailedBeforeBinding,
191 RollbackFailed,
192 FailedAfterBinding(PgPendingDatastoreCreation),
193}
194
195#[derive(Debug, Clone, Copy, PartialEq, Eq)]
197pub struct PgDatastoreReconciliationError;
198
199pub fn find_existing_datastore_binding(
202 connection: &PgDatastoreCreationConnection,
203 administrator_password: &SecretMaterial,
204 target_database: &str,
205) -> Result<Option<DatastoreIdentityBinding>, PgDatastoreReconciliationError> {
206 administrator_password.expose(|bytes| {
207 let password = std::str::from_utf8(bytes).map_err(|_| PgDatastoreReconciliationError)?;
208 let mut administrator = connection
209 .adapter(&connection.bootstrap_database, password)
210 .connect()
211 .map_err(|_| PgDatastoreReconciliationError)?;
212 let limit = administrator
213 .client()
214 .query_one(
215 "SELECT current_setting('max_identifier_length')::integer",
216 &[],
217 )
218 .map_err(|_| PgDatastoreReconciliationError)?
219 .get::<_, i32>(0);
220 if limit <= 0 || !identifiers_fit([target_database], limit as usize) {
221 return Err(PgDatastoreReconciliationError);
222 }
223 let exists = administrator
224 .client()
225 .query_one(
226 "SELECT EXISTS (SELECT 1 FROM pg_database WHERE datname = $1)",
227 &[&target_database],
228 )
229 .map_err(|_| PgDatastoreReconciliationError)?
230 .get::<_, bool>(0);
231 if !exists {
232 return Ok(None);
233 }
234 let mut target = connection
235 .adapter(target_database, password)
236 .connect()
237 .map_err(|_| PgDatastoreReconciliationError)?;
238 let bindings = read_datastore_identity_bindings(&mut target)
239 .map_err(|_| PgDatastoreReconciliationError)?;
240 let [binding] = bindings.as_slice() else {
241 return Err(PgDatastoreReconciliationError);
242 };
243 Ok(Some(binding.clone()))
244 })
245}
246
247pub fn inspect_datastore_identity_bindings(
250 connection: &PgDatastoreCreationConnection,
251 administrator_password: &SecretMaterial,
252) -> Result<Vec<DatastoreIdentityBinding>, PgDatastoreReconciliationError> {
253 administrator_password.expose(|bytes| {
254 let password = std::str::from_utf8(bytes).map_err(|_| PgDatastoreReconciliationError)?;
255 let mut administrator = connection
256 .adapter(&connection.bootstrap_database, password)
257 .connect()
258 .map_err(|_| PgDatastoreReconciliationError)?;
259 let databases = administrator
260 .client()
261 .query(
262 "SELECT datname FROM pg_database WHERE NOT datistemplate ORDER BY datname",
263 &[],
264 )
265 .map_err(|_| PgDatastoreReconciliationError)?;
266 let mut bindings = Vec::new();
267 for row in databases {
268 let database = row.get::<_, String>(0);
269 let mut inspected = connection
270 .adapter(&database, password)
271 .connect()
272 .map_err(|_| PgDatastoreReconciliationError)?;
273 if datastore_identity_table_exists(&mut inspected)
274 .map_err(|_| PgDatastoreReconciliationError)?
275 {
276 bindings.extend(
277 read_datastore_identity_bindings(&mut inspected)
278 .map_err(|_| PgDatastoreReconciliationError)?,
279 );
280 }
281 }
282 Ok(bindings)
283 })
284}
285
286pub fn validate_existing_datastore_postgresql(
290 connection: &PgDatastoreCreationConnection,
291 administrator_password: &SecretMaterial,
292 spec: &PgDatastoreValidationSpec,
293) -> Result<(), PgDatastoreReconciliationError> {
294 administrator_password.expose(|bytes| {
295 let password = std::str::from_utf8(bytes).map_err(|_| PgDatastoreReconciliationError)?;
296 let mut administrator = connection
297 .adapter(&connection.bootstrap_database, password)
298 .connect()
299 .map_err(|_| PgDatastoreReconciliationError)?;
300 let limit = administrator
301 .client()
302 .query_one(
303 "SELECT current_setting('max_identifier_length')::integer",
304 &[],
305 )
306 .map_err(|_| PgDatastoreReconciliationError)?
307 .get::<_, i32>(0);
308 if limit <= 0 || !identifiers_fit(spec.identifiers(), limit as usize) {
309 return Err(PgDatastoreReconciliationError);
310 }
311 let databases = administrator
312 .client()
313 .query(
314 "SELECT d.datname, r.rolname, EXISTS (SELECT 1 FROM aclexplode(COALESCE(d.datacl, acldefault('d', d.datdba))) a WHERE a.grantee = 0) AS public_privilege FROM pg_database d JOIN pg_roles r ON r.oid = d.datdba WHERE d.datname IN ($1, $2) ORDER BY d.datname",
315 &[&spec.target_database, &spec.catalog_database],
316 )
317 .map_err(|_| PgDatastoreReconciliationError)?;
318 if databases.len() != 2
319 || databases.iter().any(|row| row.get::<_, bool>(2))
320 || !databases.iter().any(|row| {
321 row.get::<_, String>(0) == spec.target_database
322 && row.get::<_, String>(1) == spec.target_owner
323 })
324 || !databases.iter().any(|row| {
325 row.get::<_, String>(0) == spec.catalog_database
326 && row.get::<_, String>(1) == spec.catalog_role
327 })
328 {
329 return Err(PgDatastoreReconciliationError);
330 }
331 let roles = administrator
332 .client()
333 .query(
334 "SELECT rolname, rolcanlogin, rolsuper, rolcreatedb, rolcreaterole, rolinherit, rolreplication, rolbypassrls, EXISTS (SELECT 1 FROM pg_auth_members m WHERE m.roleid = r.oid) AS has_members, EXISTS (SELECT 1 FROM pg_auth_members m WHERE m.member = r.oid) AS is_member FROM pg_roles r WHERE rolname IN ($1, $2)",
335 &[&spec.target_owner, &spec.catalog_role],
336 )
337 .map_err(|_| PgDatastoreReconciliationError)?;
338 if roles.len() != 2
339 || roles.iter().any(|row| {
340 row.get::<_, bool>(2)
341 || row.get::<_, bool>(3)
342 || row.get::<_, bool>(4)
343 || row.get::<_, bool>(5)
344 || row.get::<_, bool>(6)
345 || row.get::<_, bool>(7)
346 })
347 || !roles.iter().any(|row| {
348 row.get::<_, String>(0) == spec.target_owner
349 && !row.get::<_, bool>(1)
350 && role_has_permitted_memberships(
351 DatastoreRoleMembershipPolicy::TargetOwner,
352 row.get::<_, bool>(8),
353 row.get::<_, bool>(9),
354 )
355 })
356 || !roles.iter().any(|row| {
357 row.get::<_, String>(0) == spec.catalog_role
358 && row.get::<_, bool>(1)
359 && role_has_permitted_memberships(
360 DatastoreRoleMembershipPolicy::IsolatedCatalog,
361 row.get::<_, bool>(8),
362 row.get::<_, bool>(9),
363 )
364 })
365 {
366 return Err(PgDatastoreReconciliationError);
367 }
368 let mut target = connection
369 .adapter(&spec.target_database, password)
370 .connect()
371 .map_err(|_| PgDatastoreReconciliationError)?;
372 let metadata_objects = target
373 .client()
374 .query_one(
375 "SELECT COUNT(*)::bigint, COALESCE(bool_and(r.rolname = $1), false), NOT EXISTS (SELECT 1 FROM pg_class c2 JOIN pg_namespace n2 ON n2.oid = c2.relnamespace CROSS JOIN LATERAL aclexplode(c2.relacl) a WHERE n2.nspname = 'public' AND a.grantee = 0), NOT EXISTS (SELECT 1 FROM pg_proc p JOIN pg_namespace n3 ON n3.oid = p.pronamespace JOIN pg_roles r3 ON r3.oid = p.proowner WHERE n3.nspname = 'public' AND r3.rolname <> $1 AND NOT (p.oid = to_regprocedure('public.tre_access_request_study_exists(uuid)') AND r3.rolname = 'tre_access_request_lookup' AND NOT r3.rolcanlogin AND NOT r3.rolsuper AND NOT r3.rolcreatedb AND NOT r3.rolcreaterole AND NOT r3.rolinherit AND NOT r3.rolreplication AND r3.rolbypassrls AND NOT EXISTS (SELECT 1 FROM pg_auth_members m WHERE m.roleid = r3.oid OR m.member = r3.oid))), NOT EXISTS (SELECT 1 FROM pg_type t JOIN pg_namespace n4 ON n4.oid = t.typnamespace JOIN pg_roles r4 ON r4.oid = t.typowner WHERE n4.nspname = 'public' AND r4.rolname <> $1) FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace JOIN pg_roles r ON r.oid = c.relowner WHERE n.nspname = 'public'",
376 &[&spec.target_owner],
377 )
378 .map_err(|_| PgDatastoreReconciliationError)?;
379 if metadata_objects.get::<_, i64>(0) == 0
380 || !metadata_objects.get::<_, bool>(1)
381 || !metadata_objects.get::<_, bool>(2)
382 || !metadata_objects.get::<_, bool>(3)
383 || !metadata_objects.get::<_, bool>(4)
384 {
385 return Err(PgDatastoreReconciliationError);
386 }
387 let mut catalog = connection
388 .adapter(&spec.catalog_database, password)
389 .connect()
390 .map_err(|_| PgDatastoreReconciliationError)?;
391 let schema = catalog
392 .client()
393 .query(
394 "SELECT n.nspname, r.rolname, EXISTS (SELECT 1 FROM aclexplode(COALESCE(n.nspacl, acldefault('n', n.nspowner))) a WHERE a.grantee = 0) AS public_privilege FROM pg_namespace n JOIN pg_roles r ON r.oid = n.nspowner WHERE n.nspname IN ($1, 'public')",
395 &[&spec.catalog_schema],
396 )
397 .map_err(|_| PgDatastoreReconciliationError)?;
398 if schema.len() != 2
399 || schema.iter().any(|row| row.get::<_, bool>(2))
400 || !schema.iter().any(|row| {
401 row.get::<_, String>(0) == spec.catalog_schema
402 && row.get::<_, String>(1) == spec.catalog_role
403 })
404 || !schema.iter().any(|row| {
405 row.get::<_, String>(0) == "public"
406 && row.get::<_, String>(1) == "pg_database_owner"
407 })
408 {
409 return Err(PgDatastoreReconciliationError);
410 }
411 Ok(())
412 })
413}
414
415#[derive(Debug, Clone, Copy, PartialEq, Eq)]
416enum DatastoreRoleMembershipPolicy {
417 TargetOwner,
418 IsolatedCatalog,
419}
420
421fn role_has_permitted_memberships(
422 policy: DatastoreRoleMembershipPolicy,
423 has_members: bool,
424 is_member: bool,
425) -> bool {
426 !is_member && (policy == DatastoreRoleMembershipPolicy::TargetOwner || !has_members)
429}
430
431pub fn begin_datastore_creation(
436 connection: &PgDatastoreCreationConnection,
437 administrator_password: &SecretMaterial,
438 catalog_password: &SecretMaterial,
439 spec: PgDatastoreCreationSpec,
440) -> PgDatastoreCreationStart {
441 administrator_password.expose(|administrator_bytes| {
442 let Ok(administrator_password) = std::str::from_utf8(administrator_bytes) else {
443 return PgDatastoreCreationStart::FailedBeforeBinding;
444 };
445 let administrator =
446 connection.adapter(&connection.bootstrap_database, administrator_password);
447 let target = connection.adapter(&spec.target_database, administrator_password);
448 let catalog = connection.adapter(&spec.catalog_database, administrator_password);
449 catalog_password.expose(|catalog_bytes| {
450 let Ok(catalog_password) = std::str::from_utf8(catalog_bytes) else {
451 return PgDatastoreCreationStart::FailedBeforeBinding;
452 };
453 begin_datastore_creation_with_adapters(
454 &administrator,
455 &target,
456 &catalog,
457 spec,
458 catalog_password,
459 )
460 })
461 })
462}
463
464fn begin_datastore_creation_with_adapters(
465 administrator: &PgMetadataAdapter,
466 target: &PgMetadataAdapter,
467 catalog: &PgMetadataAdapter,
468 spec: PgDatastoreCreationSpec,
469 catalog_password: &str,
470) -> PgDatastoreCreationStart {
471 let Ok(mut admin) = administrator.connect() else {
472 return PgDatastoreCreationStart::FailedBeforeBinding;
473 };
474 let lock_key = format!("ahri-tre:create:{}", spec.target_database);
475 let Ok(lock_row) = admin.client().query_one(
476 "SELECT pg_try_advisory_lock(hashtextextended($1, 0))",
477 &[&lock_key],
478 ) else {
479 return PgDatastoreCreationStart::FailedBeforeBinding;
480 };
481 if !lock_row.get::<_, bool>(0) {
482 return PgDatastoreCreationStart::Conflict;
483 }
484
485 let Ok(limit_row) = admin.client().query_one(
486 "SELECT current_setting('max_identifier_length')::integer",
487 &[],
488 ) else {
489 return PgDatastoreCreationStart::FailedBeforeBinding;
490 };
491 let limit = limit_row.get::<_, i32>(0);
492 if limit <= 0 || !identifiers_fit(spec.identifiers(), limit as usize) {
493 return PgDatastoreCreationStart::FailedBeforeBinding;
494 }
495
496 let Ok(database_row) = admin.client().query_one(
497 "SELECT EXISTS (SELECT 1 FROM pg_database WHERE datname IN ($1, $2))",
498 &[&spec.target_database, &spec.catalog_database],
499 ) else {
500 return PgDatastoreCreationStart::FailedBeforeBinding;
501 };
502 let Ok(role_row) = admin.client().query_one(
503 "SELECT EXISTS (SELECT 1 FROM pg_roles WHERE rolname IN ($1, $2))",
504 &[&spec.target_owner, &spec.catalog_role],
505 ) else {
506 return PgDatastoreCreationStart::FailedBeforeBinding;
507 };
508 if database_row.get::<_, bool>(0) || role_row.get::<_, bool>(0) {
509 return PgDatastoreCreationStart::Conflict;
510 }
511
512 if admin
513 .client()
514 .batch_execute(&format!(
515 "CREATE ROLE {} NOLOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOINHERIT",
516 sql_identifier(&spec.target_owner),
517 ))
518 .is_err()
519 {
520 return PgDatastoreCreationStart::FailedBeforeBinding;
521 }
522
523 if admin
524 .client()
525 .batch_execute(&format!(
526 "CREATE DATABASE {} OWNER {}",
527 sql_identifier(&spec.target_database),
528 sql_identifier(&spec.target_owner),
529 ))
530 .is_err()
531 {
532 let _ = admin
533 .client()
534 .batch_execute(&format!("DROP ROLE {}", sql_identifier(&spec.target_owner),));
535 return PgDatastoreCreationStart::FailedBeforeBinding;
536 }
537
538 let prepared_store = (|| {
539 let mut store = target.connect().map_err(|_| ())?;
540 store
541 .client()
542 .batch_execute(&format!(
543 "REVOKE ALL ON DATABASE {} FROM PUBLIC; SET ROLE {}",
544 sql_identifier(&spec.target_database),
545 sql_identifier(&spec.target_owner),
546 ))
547 .map_err(|_| ())?;
548 store.bootstrap_schema().map_err(|_| ())?;
549 insert_datastore_identity_binding(&mut store, &spec.binding).map_err(|_| ())?;
550 crate::governance::GovernanceUpgrade::prepare(&mut store, &spec.governance)
551 .map_err(|_| ())?;
552 Ok::<_, ()>(store)
553 })();
554 let Ok(store) = prepared_store else {
555 let rollback = admin.client().batch_execute(&format!(
558 "DROP DATABASE {} WITH (FORCE); DROP ROLE {}",
559 sql_identifier(&spec.target_database),
560 sql_identifier(&spec.target_owner),
561 ));
562 return if rollback.is_ok() {
563 PgDatastoreCreationStart::FailedBeforeBinding
564 } else {
565 PgDatastoreCreationStart::RollbackFailed
566 };
567 };
568
569 let mut pending = PgPendingDatastoreCreation {
570 _administrator: admin,
571 store,
572 };
573 let post_binding_result = (|| {
574 pending
575 ._administrator
576 .client()
577 .batch_execute(&format!(
578 "CREATE ROLE {} LOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOINHERIT PASSWORD {}",
579 sql_identifier(&spec.catalog_role),
580 sql_literal(catalog_password),
581 ))
582 .map_err(|_| ())?;
583 pending
584 ._administrator
585 .client()
586 .batch_execute(&format!(
587 "CREATE DATABASE {} OWNER {}",
588 sql_identifier(&spec.catalog_database),
589 sql_identifier(&spec.catalog_role),
590 ))
591 .map_err(|_| ())?;
592
593 let mut catalog = catalog.connect().map_err(|_| ())?;
594 catalog
595 .client()
596 .batch_execute(&format!(
597 "REVOKE ALL ON DATABASE {} FROM PUBLIC; CREATE SCHEMA {} AUTHORIZATION {}; REVOKE ALL ON SCHEMA public FROM PUBLIC;",
598 sql_identifier(&spec.catalog_database),
599 sql_identifier(&spec.catalog_schema),
600 sql_identifier(&spec.catalog_role),
601 ))
602 .map_err(|_| ())
603 })();
604
605 if post_binding_result.is_err() {
606 PgDatastoreCreationStart::FailedAfterBinding(pending)
607 } else {
608 PgDatastoreCreationStart::Started(pending)
609 }
610}
611
612fn identifiers_fit<'a>(identifiers: impl IntoIterator<Item = &'a str>, limit: usize) -> bool {
613 identifiers
614 .into_iter()
615 .all(|identifier| !identifier.is_empty() && identifier.len() <= limit)
616}
617
618#[derive(Debug, Clone, Copy, PartialEq, Eq)]
620pub enum PgDatastoreCredentialRotationError {
621 Restored,
623 StateUncertain,
625}
626
627#[allow(clippy::too_many_arguments)]
628pub fn apply_datastore_credential_rotation(
629 connection: &PgDatastoreCreationConnection,
630 administrator_password: &SecretMaterial,
631 prior_password: &SecretMaterial,
632 staged_password: &SecretMaterial,
633 target_database: &str,
634 catalog_database: &str,
635 catalog_role: &str,
636 binding: &DatastoreIdentityBinding,
637 new_version: i32,
638 rotated_at: chrono::DateTime<chrono::Utc>,
639) -> Result<DatastoreIdentityBinding, PgDatastoreCredentialRotationError> {
640 administrator_password.expose(|administrator_bytes| {
641 let administrator_password = std::str::from_utf8(administrator_bytes)
642 .map_err(|_| PgDatastoreCredentialRotationError::Restored)?;
643 prior_password.expose(|prior_bytes| {
644 let prior_password = std::str::from_utf8(prior_bytes)
645 .map_err(|_| PgDatastoreCredentialRotationError::Restored)?;
646 staged_password.expose(|staged_bytes| {
647 let staged_password = std::str::from_utf8(staged_bytes)
648 .map_err(|_| PgDatastoreCredentialRotationError::Restored)?;
649 let mut administrator = rotation_administrator(
650 connection,
651 administrator_password,
652 binding.datastore_id,
653 )?;
654 set_role_password(&mut administrator, catalog_role, staged_password)?;
655 let applied = verify_role_password(
656 connection,
657 catalog_database,
658 catalog_role,
659 staged_password,
660 )
661 .and_then(|()| {
662 write_rotation_binding(
663 connection,
664 administrator_password,
665 target_database,
666 binding,
667 new_version,
668 rotated_at,
669 )
670 });
671 match applied {
672 Ok(binding) => Ok(binding),
673 Err(_) => {
674 let restored_password =
675 set_role_password(&mut administrator, catalog_role, prior_password)
676 .and_then(|()| {
677 verify_role_password(
678 connection,
679 catalog_database,
680 catalog_role,
681 prior_password,
682 )
683 });
684 let restored_binding = restore_rotation_binding(
685 connection,
686 administrator_password,
687 target_database,
688 binding,
689 );
690 if restored_password.is_ok() && restored_binding.is_ok() {
691 Err(PgDatastoreCredentialRotationError::Restored)
692 } else {
693 Err(PgDatastoreCredentialRotationError::StateUncertain)
694 }
695 }
696 }
697 })
698 })
699 })
700}
701
702pub fn rollback_datastore_credential_rotation(
703 connection: &PgDatastoreCreationConnection,
704 administrator_password: &SecretMaterial,
705 prior_password: &SecretMaterial,
706 target_database: &str,
707 catalog_role: &str,
708 catalog_database: &str,
709 binding: &DatastoreIdentityBinding,
710) -> Result<(), PgDatastoreCredentialRotationError> {
711 administrator_password.expose(|administrator_bytes| {
712 let administrator_password = std::str::from_utf8(administrator_bytes)
713 .map_err(|_| PgDatastoreCredentialRotationError::StateUncertain)?;
714 prior_password.expose(|prior_bytes| {
715 let prior_password = std::str::from_utf8(prior_bytes)
716 .map_err(|_| PgDatastoreCredentialRotationError::StateUncertain)?;
717 let mut administrator =
718 rotation_administrator(connection, administrator_password, binding.datastore_id)?;
719 set_role_password(&mut administrator, catalog_role, prior_password)
720 .and_then(|()| {
721 verify_role_password(connection, catalog_database, catalog_role, prior_password)
722 })
723 .and_then(|()| {
724 restore_rotation_binding(
725 connection,
726 administrator_password,
727 target_database,
728 binding,
729 )
730 .map(|_| ())
731 })
732 .map_err(|_| PgDatastoreCredentialRotationError::StateUncertain)
733 })
734 })
735}
736
737fn rotation_administrator(
738 connection: &PgDatastoreCreationConnection,
739 administrator_password: &str,
740 datastore_id: uuid::Uuid,
741) -> Result<PgMetadataConnection, PgDatastoreCredentialRotationError> {
742 let mut administrator = connection
743 .adapter(&connection.bootstrap_database, administrator_password)
744 .connect()
745 .map_err(|_| PgDatastoreCredentialRotationError::Restored)?;
746 let key = format!("ahri-tre:credential-rotation:{datastore_id}");
747 let row = administrator
748 .client()
749 .query_one(
750 "SELECT pg_try_advisory_lock(hashtextextended($1, 0))",
751 &[&key],
752 )
753 .map_err(|_| PgDatastoreCredentialRotationError::Restored)?;
754 if row.get::<_, bool>(0) {
755 Ok(administrator)
756 } else {
757 Err(PgDatastoreCredentialRotationError::Restored)
758 }
759}
760
761fn set_role_password(
762 administrator: &mut PgMetadataConnection,
763 role: &str,
764 password: &str,
765) -> Result<(), PgDatastoreCredentialRotationError> {
766 classify_password_mutation_result(administrator.client().batch_execute(&format!(
767 "ALTER ROLE {} PASSWORD {}",
768 sql_identifier(role),
769 sql_literal(password)
770 )))
771}
772
773fn classify_password_mutation_result<T, E>(
774 result: Result<T, E>,
775) -> Result<T, PgDatastoreCredentialRotationError> {
776 result.map_err(|_| PgDatastoreCredentialRotationError::StateUncertain)
777}
778
779fn verify_role_password(
780 connection: &PgDatastoreCreationConnection,
781 database: &str,
782 role: &str,
783 password: &str,
784) -> Result<(), PgDatastoreCredentialRotationError> {
785 let mut verified = connection
786 .adapter_for(database, role, password)
787 .connect()
788 .map_err(|_| PgDatastoreCredentialRotationError::Restored)?;
789 let row = verified
790 .client()
791 .query_one("SELECT current_user", &[])
792 .map_err(|_| PgDatastoreCredentialRotationError::Restored)?;
793 if row.get::<_, String>(0) == role {
794 Ok(())
795 } else {
796 Err(PgDatastoreCredentialRotationError::Restored)
797 }
798}
799
800fn write_rotation_binding(
801 connection: &PgDatastoreCreationConnection,
802 administrator_password: &str,
803 target_database: &str,
804 binding: &DatastoreIdentityBinding,
805 version: i32,
806 rotated_at: chrono::DateTime<chrono::Utc>,
807) -> Result<DatastoreIdentityBinding, PgDatastoreCredentialRotationError> {
808 let mut expected = binding.clone();
809 expected.credential_version = version;
810 expected.credential_last_rotated_at = Some(rotated_at);
811 let reference = expected
812 .managed_secret_ref
813 .clone()
814 .ok_or(PgDatastoreCredentialRotationError::Restored)?;
815 let fingerprint = expected.computed_binding_fingerprint();
816 let mut target = connection
817 .adapter(target_database, administrator_password)
818 .connect()
819 .map_err(|_| PgDatastoreCredentialRotationError::Restored)?;
820 crate::update_datastore_identity_credential(
821 &mut target,
822 binding.datastore_id,
823 &reference,
824 version,
825 rotated_at,
826 &fingerprint,
827 )
828 .map_err(|_| PgDatastoreCredentialRotationError::Restored)
829}
830
831fn restore_rotation_binding(
832 connection: &PgDatastoreCreationConnection,
833 administrator_password: &str,
834 target_database: &str,
835 binding: &DatastoreIdentityBinding,
836) -> Result<DatastoreIdentityBinding, PgDatastoreCredentialRotationError> {
837 let reference = binding
838 .managed_secret_ref
839 .as_deref()
840 .ok_or(PgDatastoreCredentialRotationError::Restored)?;
841 let rotated_at = binding
842 .credential_last_rotated_at
843 .ok_or(PgDatastoreCredentialRotationError::Restored)?;
844 let mut target = connection
845 .adapter(target_database, administrator_password)
846 .connect()
847 .map_err(|_| PgDatastoreCredentialRotationError::Restored)?;
848 crate::update_datastore_identity_credential(
849 &mut target,
850 binding.datastore_id,
851 reference,
852 binding.credential_version,
853 rotated_at,
854 &binding.binding_fingerprint,
855 )
856 .map_err(|_| PgDatastoreCredentialRotationError::Restored)
857}
858
859fn sql_identifier(value: &str) -> String {
860 format!("\"{}\"", value.replace('"', "\"\""))
861}
862
863fn sql_literal(value: &str) -> String {
864 format!("'{}'", value.replace('\'', "''"))
865}
866
867impl PgDatastoreCreationConnection {
868 pub fn upgrade_governance(
871 &self,
872 database: &str,
873 owner: &str,
874 credential: &SecretMaterial,
875 appointment: &crate::governance::GovernanceAppointment,
876 ) -> Result<crate::governance::GovernanceMaintenance, crate::PgMetaError> {
877 credential.expose(|bytes| {
878 let password=std::str::from_utf8(bytes).map_err(|_|crate::PgMetaError::Decode { field:"administrator",value:String::new(),message:"administrator credential unavailable".into() })?;
879 let mut connection=self.adapter(database,password).connect()?;
880 let actual=connection.client().query_one("SELECT pg_get_userbyid(datdba) FROM pg_database WHERE datname=current_database()",&[]).map_err(crate::PgMetaError::Query)?.get::<_,String>(0);
881 if actual!=owner { return Err(crate::PgMetaError::Decode {field:"metadata_owner",value:String::new(),message:"metadata owner mismatch".into()}); }
882 connection.client().batch_execute(&format!("SET ROLE {}",sql_identifier(owner))).map_err(crate::PgMetaError::Query)?;
883 crate::governance::GovernanceUpgrade::prepare(&mut connection,appointment)?;
884 connection.client().batch_execute("RESET ROLE").map_err(crate::PgMetaError::Query)?;
885 Ok(crate::PgDatasetMaintenance::new(connection)?.governance())
886 })
887 }
888}
889
890#[cfg(test)]
891mod tests {
892 use super::{
893 DatastoreRoleMembershipPolicy, PgDatastoreCredentialRotationError,
894 classify_password_mutation_result, identifiers_fit, role_has_permitted_memberships,
895 };
896
897 #[test]
898 fn validates_identifiers_by_server_reported_byte_limit() {
899 assert!(identifiers_fit(["target", "catalog", "role"], 63));
900 assert!(!identifiers_fit([""], 63));
901 assert!(!identifiers_fit(["a".repeat(64).as_str()], 63));
902 assert!(!identifiers_fit(["é".repeat(32).as_str()], 63));
903 }
904
905 #[test]
906 fn transport_failure_after_password_mutation_is_state_uncertain() {
907 let result = classify_password_mutation_result::<(), _>(Err("connection lost"));
908
909 assert_eq!(
910 result,
911 Err(PgDatastoreCredentialRotationError::StateUncertain)
912 );
913 }
914
915 #[test]
916 fn datastore_owner_may_be_granted_to_an_entry_group_without_inheriting_another_role() {
917 assert!(role_has_permitted_memberships(
918 DatastoreRoleMembershipPolicy::TargetOwner,
919 true,
920 false
921 ));
922 assert!(!role_has_permitted_memberships(
923 DatastoreRoleMembershipPolicy::TargetOwner,
924 false,
925 true
926 ));
927 }
928
929 #[test]
930 fn catalog_role_must_have_no_memberships_in_either_direction() {
931 assert!(role_has_permitted_memberships(
932 DatastoreRoleMembershipPolicy::IsolatedCatalog,
933 false,
934 false
935 ));
936 assert!(!role_has_permitted_memberships(
937 DatastoreRoleMembershipPolicy::IsolatedCatalog,
938 true,
939 false
940 ));
941 assert!(!role_has_permitted_memberships(
942 DatastoreRoleMembershipPolicy::IsolatedCatalog,
943 false,
944 true
945 ));
946 }
947}