Skip to main content

ahri_tre_app/
datastore_reachability.rs

1//! Safe application workflow for probing a retained Datastore Session.
2
3use ahri_tre_types::EncryptionMode;
4use uuid::Uuid;
5
6use crate::{AppError, DataStoreSession};
7
8/// Safe PostgreSQL readiness facts returned by the application workflow.
9#[derive(Debug, Clone, PartialEq, Eq)]
10pub struct MetadataStoreReachability {
11    pub database: String,
12    pub server_version: String,
13    pub schema_bootstrapped: bool,
14}
15
16/// Safe DuckLake readiness facts returned by the application workflow.
17#[derive(Debug, Clone, PartialEq, Eq)]
18pub struct LakeReachability {
19    pub extension_version: String,
20    pub configured_encryption_mode: EncryptionMode,
21    pub detected_encryption_mode: EncryptionMode,
22    pub automatic_migration: bool,
23    pub snapshot_count: i64,
24    pub table_count: i64,
25}
26
27/// Safe result of reaching one identity-bound Datastore through a live Session.
28#[derive(Debug, Clone, PartialEq, Eq)]
29pub struct DatastoreReachability {
30    pub datastore_id: Uuid,
31    pub datastore: String,
32    pub metadata: MetadataStoreReachability,
33    pub lake: LakeReachability,
34}
35
36/// Probes PostgreSQL and DuckLake through the capabilities retained by one
37/// already-authorized Session. Configuration, Secret stores, and topology are
38/// neither reopened nor projected by this workflow.
39pub fn inspect_datastore_session(
40    session: &DataStoreSession,
41) -> Result<DatastoreReachability, AppError> {
42    let runtime = session.runtime();
43    let binding = runtime.identity_binding.as_ref().ok_or_else(|| {
44        AppError::Validation("Datastore identity binding is unavailable".to_string())
45    })?;
46
47    let metadata = (|| {
48        let metadata = if let Some(mut connection) = session.direct_store_connection() {
49            let health = connection.health_check().map_err(|_| {
50                AppError::Infrastructure("metadata-store readiness check failed".to_string())
51            })?;
52            MetadataStoreReachability {
53                database: health.database,
54                server_version: health.server_version,
55                schema_bootstrapped: health.schema_bootstrapped,
56            }
57        } else if let Some(connection) = session.oauth_store_connection() {
58            let health = connection.health_check().map_err(|_| {
59                AppError::Infrastructure("metadata-store readiness check failed".to_string())
60            })?;
61            MetadataStoreReachability {
62                database: health.database,
63                server_version: health.server_version,
64                schema_bootstrapped: health.schema_bootstrapped,
65            }
66        } else {
67            return Err(AppError::Validation(
68                "metadata-store capability is unavailable".to_string(),
69            ));
70        };
71        Ok(metadata)
72    })();
73    session.observations().record(
74        ahri_tre_protocol::diagnostics::DiagnosticComponent::Metadata,
75        ahri_tre_protocol::diagnostics::DiagnosticObservationOrigin::Workflow,
76        metadata.is_ok(),
77    );
78    let metadata = metadata?;
79
80    verify_datastore_identity(
81        &metadata.database,
82        &binding.datastore_name,
83        &runtime.profile.dbname,
84    )?;
85
86    let mut catalog_observations = ahri_tre_lake::LakeCatalogObservations::default();
87    let lake_health = ahri_tre_lake::DuckLakeAdapter::new(&runtime.lake.data_path)
88        .health_check_observed(session.lake_connection(), &mut catalog_observations);
89    session.observations().record_catalog(
90        &catalog_observations,
91        ahri_tre_protocol::diagnostics::DiagnosticObservationOrigin::Workflow,
92    );
93    session.observations().record(
94        ahri_tre_protocol::diagnostics::DiagnosticComponent::Lake,
95        ahri_tre_protocol::diagnostics::DiagnosticObservationOrigin::Workflow,
96        lake_health.is_ok(),
97    );
98    let lake_health = lake_health
99        .map_err(|_| AppError::Infrastructure("Lake readiness check failed".to_string()))?;
100    let lake = LakeReachability {
101        extension_version: lake_health.extension_version,
102        configured_encryption_mode: runtime.profile.lake.auth.encryption_mode,
103        detected_encryption_mode: lake_health.detected_encryption_mode,
104        automatic_migration: runtime.lake.automatic_migration,
105        snapshot_count: lake_health.snapshot_count,
106        table_count: lake_health.table_count,
107    };
108
109    Ok(DatastoreReachability {
110        datastore_id: binding.datastore_id,
111        datastore: binding.datastore_name.clone(),
112        metadata,
113        lake,
114    })
115}
116
117fn verify_datastore_identity(
118    metadata_database: &str,
119    binding_datastore: &str,
120    profile_database: &str,
121) -> Result<(), AppError> {
122    if metadata_database == binding_datastore && binding_datastore == profile_database {
123        return Ok(());
124    }
125    Err(AppError::Validation(
126        "Datastore identity does not match the retained Session".to_string(),
127    ))
128}
129
130#[cfg(test)]
131mod tests {
132    use super::*;
133
134    #[test]
135    fn datastore_identity_must_match_the_live_store_binding_and_selected_profile() {
136        assert!(verify_datastore_identity("expected", "expected", "expected").is_ok());
137
138        for identities in [
139            ("other", "expected", "expected"),
140            ("expected", "other", "expected"),
141            ("expected", "expected", "other"),
142        ] {
143            let error = verify_datastore_identity(identities.0, identities.1, identities.2)
144                .expect_err("a mismatched Datastore identity must fail closed");
145            assert!(matches!(
146                error,
147                AppError::Validation(message)
148                    if message == "Datastore identity does not match the retained Session"
149            ));
150        }
151    }
152}