ahri_tre_lake/
catalog_observations.rs1use 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#[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
116fn 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 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#[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}