use std::path::Path;
use std::sync::Arc;
use pgevolve_core::catalog::cluster::read_cluster_catalog;
use pgevolve_core::diff::cluster::{ClusterChangeSet, diff_cluster};
use pgevolve_core::ir::cluster::catalog::ClusterCatalog;
use pgevolve_core::lint::Finding;
use pgevolve_core::lint::universal::check_cluster_changeset;
use pgevolve_core::parse::cluster::parse_cluster_directory;
use pgevolve_core::plan::cluster_rewrite::emit_cluster_changes;
use pgevolve_core::plan::raw_step::RawStep;
use crate::cluster_config::ClusterConfig;
use crate::pg_querier::PgCatalogQuerier;
#[derive(Debug)]
pub struct ClusterPlan {
pub steps: Vec<RawStep>,
pub source: ClusterCatalog,
pub target: ClusterCatalog,
pub changes: ClusterChangeSet,
pub advisory_findings: Vec<Finding>,
}
#[derive(Debug, thiserror::Error)]
pub enum ClusterPlanError {
#[error("parse error: {0}")]
Parse(#[from] pgevolve_core::parse::error::ParseError),
#[error("catalog read error: {0}")]
Catalog(#[from] pgevolve_core::catalog::error::CatalogError),
#[error("connection error: {0}")]
Connection(String),
}
pub async fn build_cluster_plan(
project_root: &Path,
cfg: &ClusterConfig,
) -> Result<ClusterPlan, ClusterPlanError> {
let roles_dir = project_root.join("roles");
let source = parse_cluster_directory(&roles_dir)?;
let (client, connection) = tokio_postgres::connect(&cfg.connection.dsn, tokio_postgres::NoTls)
.await
.map_err(|e| ClusterPlanError::Connection(e.to_string()))?;
tokio::spawn(async move {
if let Err(err) = connection.await {
tracing::debug!(?err, "cluster catalog connection task ended");
}
});
let bootstrap_roles: Vec<String> = cfg.bootstrap.roles.clone();
let querier = PgCatalogQuerier::from_arc(Arc::new(client))
.map_err(|e| ClusterPlanError::Connection(e.to_string()))?;
let target =
tokio::task::spawn_blocking(move || read_cluster_catalog(&querier, &bootstrap_roles))
.await
.map_err(|e| ClusterPlanError::Connection(format!("join error: {e}")))?
.map_err(ClusterPlanError::Catalog)?;
let changes = diff_cluster(&target, &source);
let advisory_findings = check_cluster_changeset(&source, &changes);
let steps = emit_cluster_changes(&changes);
Ok(ClusterPlan {
steps,
source,
target,
changes,
advisory_findings,
})
}