Skip to main content

ahri_tre_pgmeta/
discovery.rs

1use crate::{
2    DirectPostgresConnectConfig, MetadataExecutor, PgMetaError, PgMetadataAdapter,
3    read_datastore_identity_bindings,
4};
5use ahri_tre_observability::{CorrelationContext, FailureCategory, Outcome, Stage};
6use ahri_tre_types::DatastoreIdentityBinding;
7
8#[derive(Debug, Clone, PartialEq, Eq)]
9pub struct DatastoreServerDiscovery {
10    pub server: String,
11    pub port: u16,
12    pub maintenance_database: String,
13    pub datastores: Vec<DatastoreCandidateInspection>,
14    pub warnings: Vec<DatastoreDiscoveryWarning>,
15}
16
17#[derive(Debug, Clone, PartialEq, Eq)]
18pub struct DatastoreCandidateInspection {
19    pub database: String,
20    pub status: DatastoreCandidateStatus,
21    pub schema: Option<crate::DatastoreSchemaStatus>,
22    pub bindings: Vec<DatastoreIdentityBinding>,
23    pub message: Option<String>,
24}
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq)]
27pub enum DatastoreCandidateStatus {
28    Ready,
29    Legacy,
30    Creating,
31    Failed,
32    Retired,
33    Pending,
34    Unsupported,
35    Dirty,
36    TooNew,
37    MissingHistory,
38    Inaccessible,
39}
40
41#[derive(Debug, Clone, PartialEq, Eq)]
42pub struct DatastoreDiscoveryWarning {
43    pub database: Option<String>,
44    pub message: String,
45}
46
47pub fn discover_datastores(
48    config: DirectPostgresConnectConfig,
49    include_system_databases: bool,
50) -> Result<DatastoreServerDiscovery, PgMetaError> {
51    discover_datastores_with_ca(config, Vec::new(), include_system_databases)
52}
53
54/// Discovers datastores through explicit structured settings and an optional
55/// configuration-projected CA bundle.
56pub fn discover_datastores_with_ca(
57    config: DirectPostgresConnectConfig,
58    root_certificates: Vec<String>,
59    include_system_databases: bool,
60) -> Result<DatastoreServerDiscovery, PgMetaError> {
61    let mut maintenance = adapter(config.clone(), &root_certificates).connect()?;
62    let database_names = list_candidate_database_names(&mut maintenance, include_system_databases)?;
63    let mut datastores = Vec::new();
64    let mut warnings = Vec::new();
65
66    for database in database_names {
67        let mut candidate_config = config.clone();
68        candidate_config.dbname.clone_from(&database);
69        record_candidate_inspection_outcome(
70            &mut datastores,
71            &mut warnings,
72            database,
73            candidate_inspection_outcome(inspect_candidate_with_ca(
74                candidate_config,
75                &root_certificates,
76            )),
77        );
78    }
79
80    Ok(DatastoreServerDiscovery {
81        server: config.host,
82        port: config.port,
83        maintenance_database: config.dbname,
84        datastores,
85        warnings,
86    })
87}
88
89/// Inspects exactly one configured datastore without enumerating or opening
90/// any other database on the PostgreSQL server.
91pub fn discover_configured_datastore_with_ca(
92    config: DirectPostgresConnectConfig,
93    database: String,
94    root_certificates: Vec<String>,
95) -> Result<DatastoreServerDiscovery, PgMetaError> {
96    discover_configured_datastore_with_ca_and_routing_address(
97        config,
98        database,
99        None,
100        root_certificates,
101        None,
102    )
103}
104
105/// Inspects one configured datastore while keeping transport routing distinct
106/// from the logical PostgreSQL host used for TLS verification.
107pub fn discover_configured_datastore_with_ca_and_routing_address(
108    mut config: DirectPostgresConnectConfig,
109    database: String,
110    routing_address: Option<&str>,
111    root_certificates: Vec<String>,
112    context: Option<&CorrelationContext>,
113) -> Result<DatastoreServerDiscovery, PgMetaError> {
114    let span = context.map(|context| context.span(Stage::Metadata));
115    let server = config.host.clone();
116    let port = config.port;
117    config.dbname.clone_from(&database);
118    let mut datastores = Vec::new();
119    let mut warnings = Vec::new();
120    let inspection = inspect_candidate_with_adapter(
121        database.clone(),
122        PgMetadataAdapter::direct_with_ca_and_routing_address(
123            config,
124            routing_address,
125            root_certificates,
126        ),
127    );
128    if let Some(span) = span {
129        let category = match &inspection {
130            Err(error) => Some(error.operational_category()),
131            Ok(Some(candidate))
132                if matches!(
133                    candidate.status,
134                    DatastoreCandidateStatus::Legacy
135                        | DatastoreCandidateStatus::Pending
136                        | DatastoreCandidateStatus::Unsupported
137                        | DatastoreCandidateStatus::Dirty
138                        | DatastoreCandidateStatus::TooNew
139                        | DatastoreCandidateStatus::MissingHistory
140                ) =>
141            {
142                Some(FailureCategory::Compatibility)
143            }
144            _ => None,
145        };
146        span.finish(
147            match category {
148                None => Outcome::Success,
149                Some(FailureCategory::Authentication | FailureCategory::Authorization) => {
150                    Outcome::Rejected
151                }
152                Some(_) => Outcome::Unavailable,
153            },
154            category,
155        );
156    }
157    record_candidate_inspection_outcome(
158        &mut datastores,
159        &mut warnings,
160        database.clone(),
161        candidate_inspection_outcome(inspection),
162    );
163    Ok(DatastoreServerDiscovery {
164        server,
165        port,
166        maintenance_database: database,
167        datastores,
168        warnings,
169    })
170}
171
172fn adapter(config: DirectPostgresConnectConfig, root_certificates: &[String]) -> PgMetadataAdapter {
173    if root_certificates.is_empty() {
174        PgMetadataAdapter::direct(config)
175    } else {
176        PgMetadataAdapter::direct_with_ca(config, root_certificates.to_vec())
177    }
178}
179
180enum CandidateInspectionOutcome {
181    Candidate(Box<Option<DatastoreCandidateInspection>>),
182    Omitted,
183    Inaccessible { message: String },
184}
185
186fn candidate_inspection_outcome(
187    result: Result<Option<DatastoreCandidateInspection>, PgMetaError>,
188) -> CandidateInspectionOutcome {
189    match result {
190        Ok(candidate) => CandidateInspectionOutcome::Candidate(Box::new(candidate)),
191        Err(error) if error.is_authorization_failure() => CandidateInspectionOutcome::Omitted,
192        Err(error) => CandidateInspectionOutcome::Inaccessible {
193            message: error.to_string(),
194        },
195    }
196}
197
198fn inspect_candidate_with_ca(
199    config: DirectPostgresConnectConfig,
200    root_certificates: &[String],
201) -> Result<Option<DatastoreCandidateInspection>, PgMetaError> {
202    let database = config.dbname.clone();
203    let mut connection = adapter(config, root_certificates).connect()?;
204    inspect_candidate_connection(database, &mut connection)
205}
206
207fn inspect_candidate_with_adapter(
208    database: String,
209    adapter: PgMetadataAdapter,
210) -> Result<Option<DatastoreCandidateInspection>, PgMetaError> {
211    let mut connection = adapter.connect()?;
212    inspect_candidate_connection(database, &mut connection)
213}
214
215fn record_candidate_inspection_outcome(
216    datastores: &mut Vec<DatastoreCandidateInspection>,
217    warnings: &mut Vec<DatastoreDiscoveryWarning>,
218    database: String,
219    outcome: CandidateInspectionOutcome,
220) {
221    match outcome {
222        CandidateInspectionOutcome::Candidate(candidate) => {
223            if let Some(candidate) = *candidate {
224                datastores.push(candidate);
225            }
226        }
227        CandidateInspectionOutcome::Omitted => {}
228        CandidateInspectionOutcome::Inaccessible { message } => {
229            warnings.push(DatastoreDiscoveryWarning {
230                database: Some(database.clone()),
231                message: "candidate database could not be inspected".to_string(),
232            });
233            datastores.push(DatastoreCandidateInspection {
234                database,
235                status: DatastoreCandidateStatus::Inaccessible,
236                schema: None,
237                bindings: Vec::new(),
238                message: Some(message),
239            });
240        }
241    }
242}
243
244fn list_candidate_database_names<E>(
245    executor: &mut E,
246    include_system_databases: bool,
247) -> Result<Vec<String>, PgMetaError>
248where
249    E: MetadataExecutor,
250    E::Row: DatabaseNameRow,
251{
252    let rows = if include_system_databases {
253        executor.query_many(
254            "list PostgreSQL databases for datastore discovery",
255            "SELECT datname
256               FROM pg_database
257              WHERE datallowconn
258              ORDER BY datname",
259            &[],
260        )?
261    } else {
262        executor.query_many(
263            "list PostgreSQL user databases for datastore discovery",
264            "SELECT datname
265               FROM pg_database
266              WHERE datallowconn
267                AND NOT datistemplate
268                AND datname NOT IN ('postgres', 'template0', 'template1')
269              ORDER BY datname",
270            &[],
271        )?
272    };
273    Ok(rows.iter().map(DatabaseNameRow::database_name).collect())
274}
275
276fn inspect_candidate_connection(
277    database: String,
278    connection: &mut crate::PgMetadataConnection,
279) -> Result<Option<DatastoreCandidateInspection>, PgMetaError> {
280    let signals = read_candidate_signals(connection)?;
281    if !signals.has_any_signal() {
282        return Ok(None);
283    }
284
285    let schema = crate::datastore_schema_status(connection)?;
286    let bindings = if signals.identity_table {
287        read_datastore_identity_bindings(connection)?
288    } else {
289        Vec::new()
290    };
291    let status = classify_candidate(Some(&schema), &bindings, signals.identity_table);
292
293    Ok(Some(DatastoreCandidateInspection {
294        database,
295        status,
296        schema: Some(schema),
297        bindings,
298        message: None,
299    }))
300}
301
302#[derive(Debug, Clone, Copy, PartialEq, Eq)]
303struct CandidateSignals {
304    identity_table: bool,
305    migration_history: bool,
306    metadata_table_count: i64,
307}
308
309impl CandidateSignals {
310    const fn has_any_signal(self) -> bool {
311        self.identity_table || self.migration_history || self.metadata_table_count > 0
312    }
313}
314
315fn read_candidate_signals<E>(executor: &mut E) -> Result<CandidateSignals, PgMetaError>
316where
317    E: MetadataExecutor,
318    E::Row: CandidateSignalsRow,
319{
320    let metadata_tables = crate::METADATA_TABLES
321        .iter()
322        .map(|table| format!("'{table}'"))
323        .collect::<Vec<_>>()
324        .join(", ");
325    let statement = format!(
326        "SELECT
327            to_regclass('public.datastore_identity') IS NOT NULL AS identity_table,
328            to_regclass('public.ahri_tre_schema_migrations') IS NOT NULL AS migration_history,
329            (
330                SELECT count(*)::bigint
331                  FROM information_schema.tables
332                 WHERE table_schema = 'public'
333                   AND table_name IN ({metadata_tables})
334            ) AS metadata_table_count"
335    );
336    let row = executor.query_one("inspect datastore discovery signals", &statement, &[])?;
337    Ok(CandidateSignals {
338        identity_table: row.identity_table(),
339        migration_history: row.migration_history(),
340        metadata_table_count: row.metadata_table_count(),
341    })
342}
343
344fn classify_candidate(
345    schema: Option<&crate::DatastoreSchemaStatus>,
346    bindings: &[DatastoreIdentityBinding],
347    identity_table: bool,
348) -> DatastoreCandidateStatus {
349    if bindings
350        .iter()
351        .any(|binding| binding.lifecycle_state == ahri_tre_types::DatastoreLifecycleState::Ready)
352    {
353        return if schema
354            .map(|schema| {
355                matches!(
356                    schema.status,
357                    crate::DatastoreSchemaCompatibility::Current
358                        | crate::DatastoreSchemaCompatibility::Pending
359                )
360            })
361            .unwrap_or(false)
362        {
363            DatastoreCandidateStatus::Ready
364        } else {
365            DatastoreCandidateStatus::Pending
366        };
367    }
368
369    if let Some(binding) = bindings.first() {
370        return match binding.lifecycle_state {
371            ahri_tre_types::DatastoreLifecycleState::Creating => DatastoreCandidateStatus::Creating,
372            ahri_tre_types::DatastoreLifecycleState::Failed => DatastoreCandidateStatus::Failed,
373            ahri_tre_types::DatastoreLifecycleState::Retired => DatastoreCandidateStatus::Retired,
374            ahri_tre_types::DatastoreLifecycleState::Ready => DatastoreCandidateStatus::Pending,
375        };
376    }
377
378    match schema.map(|schema| schema.status) {
379        Some(crate::DatastoreSchemaCompatibility::Current) if !identity_table => {
380            DatastoreCandidateStatus::Unsupported
381        }
382        Some(crate::DatastoreSchemaCompatibility::Pending) => DatastoreCandidateStatus::Pending,
383        Some(crate::DatastoreSchemaCompatibility::Unsupported) => {
384            DatastoreCandidateStatus::Unsupported
385        }
386        Some(crate::DatastoreSchemaCompatibility::Dirty) => DatastoreCandidateStatus::Dirty,
387        Some(crate::DatastoreSchemaCompatibility::TooNew) => DatastoreCandidateStatus::TooNew,
388        Some(crate::DatastoreSchemaCompatibility::MissingHistory) => {
389            DatastoreCandidateStatus::MissingHistory
390        }
391        Some(crate::DatastoreSchemaCompatibility::Current) => DatastoreCandidateStatus::Pending,
392        None => DatastoreCandidateStatus::Unsupported,
393    }
394}
395
396trait DatabaseNameRow {
397    fn database_name(&self) -> String;
398}
399
400impl DatabaseNameRow for postgres::Row {
401    fn database_name(&self) -> String {
402        self.get("datname")
403    }
404}
405
406trait CandidateSignalsRow {
407    fn identity_table(&self) -> bool;
408    fn migration_history(&self) -> bool;
409    fn metadata_table_count(&self) -> i64;
410}
411
412impl CandidateSignalsRow for postgres::Row {
413    fn identity_table(&self) -> bool {
414        self.get("identity_table")
415    }
416
417    fn migration_history(&self) -> bool {
418        self.get("migration_history")
419    }
420
421    fn metadata_table_count(&self) -> i64 {
422        self.get("metadata_table_count")
423    }
424}
425
426#[cfg(test)]
427mod tests {
428    use super::*;
429    use crate::error::is_authorization_failure_sqlstate;
430    use ahri_tre_types::{
431        DatastoreIdentityBinding, DatastoreLifecycleState, LakeCatalogCredentialMode,
432    };
433    use chrono::Utc;
434    use uuid::Uuid;
435
436    #[test]
437    fn classification_reports_ready_binding_as_ready_when_schema_current() {
438        let schema = schema(crate::DatastoreSchemaCompatibility::Current);
439        let bindings = vec![binding(DatastoreLifecycleState::Ready)];
440
441        assert_eq!(
442            classify_candidate(Some(&schema), &bindings, true),
443            DatastoreCandidateStatus::Ready
444        );
445    }
446
447    #[test]
448    fn classification_rejects_current_schema_without_identity_binding() {
449        let schema = schema(crate::DatastoreSchemaCompatibility::Current);
450
451        assert_eq!(
452            classify_candidate(Some(&schema), &[], false),
453            DatastoreCandidateStatus::Unsupported
454        );
455    }
456
457    #[test]
458    fn classification_reports_failed_binding() {
459        let schema = schema(crate::DatastoreSchemaCompatibility::Current);
460        let bindings = vec![binding(DatastoreLifecycleState::Failed)];
461
462        assert_eq!(
463            classify_candidate(Some(&schema), &bindings, true),
464            DatastoreCandidateStatus::Failed
465        );
466    }
467
468    #[test]
469    fn database_permission_denial_is_outside_browser_discovery_scope() {
470        assert!(is_authorization_failure_sqlstate(Some("42501")));
471
472        assert!(!is_authorization_failure_sqlstate(Some("08006")));
473        assert!(!is_authorization_failure_sqlstate(Some("28000")));
474        assert!(!is_authorization_failure_sqlstate(Some("28P01")));
475        assert!(!is_authorization_failure_sqlstate(Some("42P01")));
476        assert!(!is_authorization_failure_sqlstate(None));
477    }
478
479    #[test]
480    fn candidate_inspection_outcomes_omit_browser_denials_and_retain_other_failures() {
481        let mut datastores = Vec::new();
482        let mut warnings = Vec::new();
483
484        record_candidate_inspection_outcome(
485            &mut datastores,
486            &mut warnings,
487            "browser-denied".to_string(),
488            CandidateInspectionOutcome::Omitted,
489        );
490        record_candidate_inspection_outcome(
491            &mut datastores,
492            &mut warnings,
493            "network-failed".to_string(),
494            CandidateInspectionOutcome::Inaccessible {
495                message: "connection refused".to_string(),
496            },
497        );
498        record_candidate_inspection_outcome(
499            &mut datastores,
500            &mut warnings,
501            "ready".to_string(),
502            CandidateInspectionOutcome::Candidate(Box::new(Some(candidate(
503                "ready",
504                DatastoreCandidateStatus::Ready,
505            )))),
506        );
507
508        assert_eq!(
509            datastores
510                .iter()
511                .map(|candidate| (candidate.database.as_str(), candidate.status))
512                .collect::<Vec<_>>(),
513            vec![
514                ("network-failed", DatastoreCandidateStatus::Inaccessible),
515                ("ready", DatastoreCandidateStatus::Ready),
516            ]
517        );
518        assert_eq!(warnings.len(), 1);
519        assert_eq!(warnings[0].database.as_deref(), Some("network-failed"));
520        assert_eq!(
521            warnings[0].message,
522            "candidate database could not be inspected"
523        );
524    }
525
526    fn schema(status: crate::DatastoreSchemaCompatibility) -> crate::DatastoreSchemaStatus {
527        crate::DatastoreSchemaStatus {
528            current_version: Some(crate::CURRENT_DATASTORE_SCHEMA_VERSION.to_string()),
529            target_version: crate::CURRENT_DATASTORE_SCHEMA_VERSION.to_string(),
530            status,
531            migration_history_present: true,
532            applied_versions: Vec::new(),
533            missing_tables: Vec::new(),
534            database: "tre".to_string(),
535            user: "postgres".to_string(),
536            server_version: "18".to_string(),
537        }
538    }
539
540    fn binding(lifecycle_state: DatastoreLifecycleState) -> DatastoreIdentityBinding {
541        DatastoreIdentityBinding {
542            datastore_id: Uuid::nil(),
543            datastore_name: "tre".to_string(),
544            ducklake_catalog_database: "tre_catalog".to_string(),
545            ducklake_catalog_schema: Some("ducklake_catalog".to_string()),
546            lake_path: "/srv/tre".to_string(),
547            storage_policy_id: None,
548            lake_location: None,
549            ducklake_encryption: true,
550            lake_catalog_credential_mode: LakeCatalogCredentialMode::ManagedLocal,
551            lake_catalog_role_name: Some("tre_lake".to_string()),
552            managed_secret_ref: Some("local:tre".to_string()),
553            credential_version: 1,
554            credential_last_rotated_at: None,
555            created_at: Utc::now(),
556            creating_tool_version: None,
557            binding_fingerprint: "0".repeat(64),
558            lifecycle_state,
559            failure_code: None,
560            failure_summary: None,
561        }
562    }
563
564    fn candidate(database: &str, status: DatastoreCandidateStatus) -> DatastoreCandidateInspection {
565        DatastoreCandidateInspection {
566            database: database.to_string(),
567            status,
568            schema: None,
569            bindings: Vec::new(),
570            message: None,
571        }
572    }
573}