mod common;
use anyhow::Result;
use pgevolve_core::ir::catalog::Catalog;
use pgevolve_core::parse::parse_directory;
use pgevolve_testkit::ephemeral_pg::{EphemeralPostgres, default_pg_version, docker_available};
use common::{apply_diff, assert_convergent, connect_and_bootstrap, introspect, schemas_of};
const SOURCE_SQL: &str = "\
-- @pgevolve schema=app
CREATE SCHEMA app;
CREATE FUNCTION app.sum_sfunc(bigint, integer) RETURNS bigint \
LANGUAGE plpgsql AS $$ BEGIN RETURN $1 + $2; END $$;
CREATE FUNCTION app.sum_ffunc(bigint) RETURNS numeric \
LANGUAGE plpgsql AS $$ BEGIN RETURN $1::numeric; END $$;
CREATE AGGREGATE app.my_sum(integer) ( \
SFUNC = app.sum_sfunc, \
STYPE = bigint, \
FINALFUNC = app.sum_ffunc, \
INITCOND = '0' \
);
";
fn source_catalog() -> Result<Catalog> {
let dir = tempfile::tempdir()?;
std::fs::write(dir.path().join("0001-aggregate.sql"), SOURCE_SQL)?;
let catalog = parse_directory(dir.path(), &[])?;
Ok(catalog)
}
#[ignore = "e2e test — requires Docker; run via `cargo test -- --ignored`"]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn aggregate_round_trips_against_ephemeral_pg() {
if !docker_available() {
eprintln!("skipping: docker unavailable");
return;
}
run().await.expect("aggregate round-trip");
}
async fn run() -> Result<()> {
let source = source_catalog()?;
let managed = schemas_of(&source);
assert_eq!(
source.aggregates.len(),
1,
"source catalog should declare exactly one aggregate"
);
let pg = EphemeralPostgres::start(default_pg_version()).await?;
let mut client = connect_and_bootstrap(&pg).await?;
let outcome = apply_diff(&mut client, &Catalog::empty(), &source, &managed, None).await?;
outcome.map_err(|e| anyhow::anyhow!("apply failed: {e}"))?;
let live = introspect(&pg, &managed).await?;
assert_eq!(
live.aggregates.len(),
1,
"live database should report exactly one aggregate after apply"
);
assert_convergent(&live, &source)
}