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_sources;
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),
}
impl ClusterPlan {
pub fn to_plan(
self,
plan_id: pgevolve_core::plan::PlanId,
target_identity: String,
) -> Result<pgevolve_core::plan::Plan, pgevolve_core::plan::PlanError> {
let groups = pgevolve_core::plan::group_steps(self.steps);
pgevolve_core::plan::Plan::from_grouped_with_id(
groups,
plan_id,
target_identity,
None, pgevolve_core::VERSION,
pgevolve_core::plan::PlannerPolicy::default().planner_ruleset_version,
)
}
}
pub async fn build_cluster_plan(
project_root: &Path,
cfg: &ClusterConfig,
) -> Result<ClusterPlan, ClusterPlanError> {
let roles_dir = project_root.join("roles");
let tablespaces_dir = cfg.tablespaces.dir.clone().map_or_else(
|| project_root.join("tablespaces"),
|d| project_root.join(d),
);
let source = parse_cluster_sources(&roles_dir, &tablespaces_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, &target, &changes);
let steps = emit_cluster_changes(&changes);
Ok(ClusterPlan {
steps,
source,
target,
changes,
advisory_findings,
})
}
#[cfg(test)]
mod tests {
use super::*;
use pgevolve_core::plan::raw_step::{RawStep, StepKind, TransactionConstraint};
fn synthetic_create_role(name: &str) -> RawStep {
RawStep {
step_no: 0,
kind: StepKind::CreateRole,
destructive: false,
destructive_reason: None,
intent_id: None,
targets: vec![],
sql: format!("CREATE ROLE {name};"),
transactional: TransactionConstraint::InTransaction,
}
}
fn synthetic_drop_role(name: &str) -> RawStep {
RawStep {
step_no: 0,
kind: StepKind::DropRole,
destructive: true,
destructive_reason: Some(format!("drops role {name} (may orphan objects)")),
intent_id: None,
targets: vec![],
sql: format!("DROP ROLE {name};"),
transactional: TransactionConstraint::InTransaction,
}
}
fn empty_changes() -> pgevolve_core::diff::cluster::ClusterChangeSet {
pgevolve_core::diff::cluster::ClusterChangeSet::default()
}
fn empty_plan_id() -> pgevolve_core::plan::PlanId {
pgevolve_core::plan::PlanId::compute(
&pgevolve_core::ir::catalog::Catalog::empty(),
&pgevolve_core::ir::catalog::Catalog::empty(),
"0.0.0-test",
0,
)
.expect("compute placeholder id")
}
#[test]
fn to_plan_assigns_step_numbers_and_intents() {
let plan = ClusterPlan {
steps: vec![synthetic_create_role("a"), synthetic_drop_role("b")],
source: ClusterCatalog::empty(),
target: ClusterCatalog::empty(),
changes: empty_changes(),
advisory_findings: vec![],
};
let core_plan = plan
.to_plan(empty_plan_id(), "cluster:0000000000003039".into())
.expect("to_plan ok");
assert_eq!(core_plan.groups.len(), 1);
assert_eq!(core_plan.groups[0].steps.len(), 2);
assert_eq!(core_plan.groups[0].steps[0].step_no, 1);
assert_eq!(core_plan.groups[0].steps[1].step_no, 2);
assert_eq!(core_plan.intents.len(), 1);
assert_eq!(core_plan.intents[0].step, 2);
assert!(!core_plan.intents[0].approved);
assert_eq!(
core_plan.metadata.target_identity,
"cluster:0000000000003039"
);
}
#[test]
fn to_plan_empty_steps_produces_empty_plan() {
let plan = ClusterPlan {
steps: vec![],
source: ClusterCatalog::empty(),
target: ClusterCatalog::empty(),
changes: empty_changes(),
advisory_findings: vec![],
};
let core_plan = plan
.to_plan(empty_plan_id(), "cluster:empty".into())
.expect("empty to_plan ok");
assert!(core_plan.groups.is_empty());
assert!(core_plan.intents.is_empty());
}
}