Skip to main content

ahri_tre_lake/
catalog_observations.rs

1//! Safe evidence captured by the existing catalog attach and health read.
2use ahri_tre_observability::{CorrelationContext, FailureCategory, Outcome, Stage};
3use chrono::{DateTime, Utc};
4
5use crate::{DuckLakeAdapter, DuckLakeHealth, LakeError};
6
7#[derive(Debug, Clone, Copy, PartialEq, Eq)]
8pub enum CatalogObservationResult {
9    Passed,
10    AuthenticationRejected,
11    AuthenticationNotEstablished,
12    ReadFailed,
13    ReadNotAttempted,
14}
15
16#[derive(Debug, Clone, Copy)]
17pub struct CatalogObservation {
18    pub observed_at_utc: DateTime<Utc>,
19    pub result: CatalogObservationResult,
20}
21
22/// Operation-local evidence only. No handles, credentials, driver messages or
23/// persistent history. An unattempted stage has no fabricated success.
24#[derive(Default)]
25pub struct LakeCatalogObservations {
26    context: Option<CorrelationContext>,
27    authentication: Option<CatalogObservation>,
28    read: Option<CatalogObservation>,
29}
30
31impl LakeCatalogObservations {
32    pub fn new(context: Option<&CorrelationContext>) -> Self {
33        Self {
34            context: context.copied(),
35            ..Self::default()
36        }
37    }
38
39    pub fn authentication(&self) -> Option<CatalogObservation> {
40        self.authentication
41    }
42    pub fn read(&self) -> Option<CatalogObservation> {
43        self.read
44    }
45
46    pub(crate) fn attach(
47        &mut self,
48        connection: &duckdb::Connection,
49        sql: &str,
50        role: &str,
51    ) -> Result<(), LakeError> {
52        let span = self
53            .context
54            .map(|context| context.span(Stage::LakeCatalogAuthentication));
55        let result = connection.execute_batch(sql);
56        let (evidence, outcome, category) = match &result {
57            Ok(()) => (CatalogObservationResult::Passed, Outcome::Success, None),
58            Err(error) if authentication_rejected(error, role) => (
59                CatalogObservationResult::AuthenticationRejected,
60                Outcome::Rejected,
61                Some(FailureCategory::Authentication),
62            ),
63            Err(_) => (
64                CatalogObservationResult::AuthenticationNotEstablished,
65                Outcome::Unavailable,
66                Some(FailureCategory::Adapter),
67            ),
68        };
69        self.authentication = Some(CatalogObservation {
70            observed_at_utc: Utc::now(),
71            result: evidence,
72        });
73        if result.is_err() {
74            self.read = Some(CatalogObservation {
75                observed_at_utc: Utc::now(),
76                result: CatalogObservationResult::ReadNotAttempted,
77            });
78        }
79        if let Some(span) = span {
80            span.finish_observed(outcome, category);
81        }
82        result.map_err(LakeError::from)
83    }
84
85    pub(crate) fn health_check(
86        &mut self,
87        adapter: &DuckLakeAdapter,
88        connection: &duckdb::Connection,
89    ) -> Result<DuckLakeHealth, LakeError> {
90        let span = self
91            .context
92            .map(|context| context.span(Stage::LakeCatalogRead));
93        let result = adapter.health_check(connection);
94        self.read = Some(CatalogObservation {
95            observed_at_utc: Utc::now(),
96            result: if result.is_ok() {
97                CatalogObservationResult::Passed
98            } else {
99                CatalogObservationResult::ReadFailed
100            },
101        });
102        if let Some(span) = span {
103            span.finish_observed(
104                if result.is_ok() {
105                    Outcome::Success
106                } else {
107                    Outcome::Unavailable
108                },
109                result.as_ref().err().map(|_| FailureCategory::Query),
110            );
111        }
112        result
113    }
114}
115
116/// DuckDB's postgres extension preserves libpq's connection-error message but
117/// does not expose its SQLSTATE. Recognize only exact English FATAL messages for
118/// the already-resolved catalog role at the connection-error boundary. Unknown
119/// or localized messages stay not established. Nothing from this private string
120/// is returned, logged or retained. Query errors cannot enter this discriminator.
121fn authentication_rejected(error: &duckdb::Error, role: &str) -> bool {
122    let duckdb::Error::DuckDBFailure(_, Some(message)) = error else {
123        return false;
124    };
125    let message = message.strip_prefix("IO Error: ").unwrap_or(message);
126    // DuckLake prepends its metadata-attach context to the same preserved
127    // postgres connection error (ErrorData::Throw), without a separator.
128    let connection_error = message.starts_with("Unable to connect to Postgres at \"")
129        || (message.starts_with("Failed to attach DuckLake MetaData \"")
130            && message.contains("\"Unable to connect to Postgres at \""));
131    if !connection_error {
132        return false;
133    }
134    [
135        format!("FATAL:  role \"{role}\" is not permitted to log in"),
136        format!("FATAL:  password authentication failed for user \"{role}\""),
137    ]
138    .iter()
139    .any(|rejection| message.trim_end().ends_with(rejection))
140}
141
142/// The source-built acceptance image alone compiles this synchronization point.
143/// It adds no production switch or endpoint and never substitutes a result.
144#[cfg(feature = "acceptance-catalog-barrier")]
145pub(crate) fn after_attach_barrier() -> Result<(), LakeError> {
146    use std::{
147        fs,
148        path::Path,
149        time::{Duration, Instant},
150    };
151    let root = Path::new("/tmp/ahri-tre-catalog-read-barrier");
152    match fs::rename(root.join("armed"), root.join("attached")) {
153        Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()),
154        result => result?,
155    }
156    let start = Instant::now();
157    while !root.join("release").exists() {
158        if start.elapsed() > Duration::from_secs(25) {
159            return Err(std::io::Error::new(
160                std::io::ErrorKind::TimedOut,
161                "acceptance barrier timed out",
162            )
163            .into());
164        }
165        std::thread::sleep(Duration::from_millis(20));
166    }
167    Ok(())
168}