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
54pub 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
89pub 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
105pub 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}