Skip to main content

ahri_tre_pgmeta/
datastore_creation.rs

1//! PostgreSQL-owned mechanics for starting a configured Datastore creation saga.
2
3use 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/// Explicit PostgreSQL authority used for one configured creation attempt.
13#[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    /// Opens the configured maintenance credential behind a cleanup-only API.
40    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/// PostgreSQL names and initial binding required for a new Datastore.
81#[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/// Expected PostgreSQL ownership and privilege boundary for one retained Datastore.
126#[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
163/// PostgreSQL capabilities retained while a creation saga is in progress.
164pub 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/// Closed start result that distinguishes failures with a durable binding.
186#[derive(Debug)]
187pub enum PgDatastoreCreationStart {
188    Started(PgPendingDatastoreCreation),
189    Conflict,
190    FailedBeforeBinding,
191    RollbackFailed,
192    FailedAfterBinding(PgPendingDatastoreCreation),
193}
194
195/// Closed adapter failure for retained-state inspection and reconciliation.
196#[derive(Debug, Clone, Copy, PartialEq, Eq)]
197pub struct PgDatastoreReconciliationError;
198
199/// Finds the singleton identity binding in an existing target database without
200/// adopting an unbound or ambiguously bound database.
201pub 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
247/// Inspects every connectable database on one configured PostgreSQL server.
248/// Any database that cannot be inspected makes absence unprovable.
249pub 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
286/// Verifies retained database/schema ownership, identifier limits, role
287/// attributes, memberships, approved service-function ownership, and PUBLIC
288/// privilege revocations.
289pub 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    // A datastore owner may be granted to an entry group, but must not inherit
427    // another role. The credential-bearing catalog role stays isolated both ways.
428    !is_member && (policy == DatastoreRoleMembershipPolicy::TargetOwner || !has_members)
429}
430
431/// Creates the target database and immediately records its `creating` binding,
432/// then performs the remaining PostgreSQL side effects under the same advisory
433/// lock. Failures after the binding is durable return the pending capability so
434/// the caller can atomically mark the saga `failed`.
435pub 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        // This database was proven absent and created by this locked attempt.
556        // Removing it here is rollback, not adoption or force-reset behavior.
557        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/// Disclosure-safe failure from the PostgreSQL-owned credential authority.
619#[derive(Debug, Clone, Copy, PartialEq, Eq)]
620pub enum PgDatastoreCredentialRotationError {
621    /// The prior password and binding were authoritatively restored.
622    Restored,
623    /// PostgreSQL may require either the prior or staged password.
624    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    /// Preserving operator upgrade with the original metadata owner. Only an
869    /// already resolved administrator credential can create this capability.
870    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}