#![cfg(all(feature = "sql-postgres", feature = "migrate"))]
use std::collections::BTreeMap;
use std::sync::Arc;
use boatramp_core::deploy::DeployStore;
use boatramp_core::kv::MemoryKv;
use boatramp_core::sql::{
LedgerOrigin, MigrateDdlError, MigrationAction, MigrationStep, MigrationSubstrate,
SubstrateStepOutcome,
};
use boatramp_node::config::{ExternalDatabaseConfig, TenantIsolation, TenantScope};
use boatramp_node::managed_sql::{NodeMigrationRunner, NodeOperatorSql};
use boatramp_storage::FsStorage;
const DB: &str = "app";
const URL_ENV: &str = "BOATRAMP_MIGRATE_RUNNER_TEST_URL";
fn sql_step(id: &str, script: &str) -> MigrationStep {
MigrationStep {
id: id.to_string(),
action: MigrationAction::Sql {
script: script.to_string(),
no_transaction: false,
},
}
}
fn ext_step(id: &str, name: &str) -> MigrationStep {
MigrationStep {
id: id.to_string(),
action: MigrationAction::Extension {
name: name.to_string(),
},
}
}
fn runner_for() -> Option<NodeMigrationRunner> {
let url = std::env::var("BOATRAMP_TEST_PG_URL").ok()?;
let env = boatramp_core::env::MapEnv::new().with(URL_ENV, url);
let mut databases = BTreeMap::new();
databases.insert(
DB.to_string(),
ExternalDatabaseConfig {
kind: "postgres".to_string(),
url_env: URL_ENV.to_string(),
compute: None,
database: None,
user: None,
pool_max: Some(4),
read_only: false,
connect_timeout_secs: Some(10),
tenant: TenantIsolation::Shared,
tenant_scope: TenantScope::Project,
..Default::default()
},
);
let op = Arc::new(
NodeOperatorSql::new(
databases,
Arc::new(MemoryKv::new()),
None,
DeployStore::new(
Arc::new(FsStorage::new(
std::env::temp_dir().join("boatramp-migrate-test"),
)),
Arc::new(MemoryKv::new()),
),
)
.with_env_source(Arc::new(env)),
);
let mut allow = std::collections::BTreeSet::new();
allow.insert("citext".to_string());
Some(NodeMigrationRunner::new(op, allow))
}
async fn apply(
sub: &NodeMigrationRunner,
step: &MigrationStep,
ordinal: usize,
) -> SubstrateStepOutcome {
let eff = step.content_hash();
sub.apply_substrate_step("default", DB, step, ordinal, &eff)
.await
.expect("substrate step (infra)")
}
fn applied_ids(applied: &[boatramp_core::sql::AppliedMigration]) -> Vec<String> {
applied.iter().map(|a| a.id.clone()).collect()
}
#[tokio::test]
async fn migrate_substrate_ledger_atomicity_and_owner_ddl_on_a_real_engine() {
let Some(sub) = runner_for() else {
eprintln!("skip migrate_substrate: BOATRAMP_TEST_PG_URL unset");
return;
};
{
use boatramp_storage::sql_sqlx::{ExternalSqlKind, ExternalSqlOptions, connect};
let url = std::env::var("BOATRAMP_TEST_PG_URL").unwrap();
let c = connect(ExternalSqlKind::Postgres, &ExternalSqlOptions::new(url)).unwrap();
for stmt in [
"DROP SCHEMA IF EXISTS boatramp_migrations CASCADE",
"DROP TABLE IF EXISTS widget",
"DROP TABLE IF EXISTS gadget",
"DROP TABLE IF EXISTS baselined",
] {
let _ = c.run_script(stmt).await;
}
}
let applied = sub.preflight("default", DB).await.unwrap();
assert!(applied.is_empty(), "empty ledger");
let s1 = sql_step(
"0001_widget",
"CREATE TABLE widget (id int primary key, name text)",
);
let s2 = sql_step("0002_seed", "INSERT INTO widget (id, name) VALUES (1, 'a')");
assert!(matches!(
apply(&sub, &s1, 0).await,
SubstrateStepOutcome::Applied
));
assert!(matches!(
apply(&sub, &s2, 1).await,
SubstrateStepOutcome::Applied
));
let applied = sub.preflight("default", DB).await.unwrap();
assert_eq!(applied_ids(&applied), vec!["0001_widget", "0002_seed"]);
assert!(
applied.iter().all(|a| a.origin == "apply"),
"applied rows carry origin=apply"
);
assert_eq!(
applied[0].content_hash,
s1.content_hash(),
"the recorded hash is the effective hash the orchestrator supplied"
);
let bad = sql_step(
"0003_atomic",
"CREATE TABLE gadget (id int primary key); INSERT INTO gadget (id) VALUES ('not-an-int')",
);
assert!(
matches!(apply(&sub, &bad, 2).await, SubstrateStepOutcome::Failed(_)),
"a mid-script error is a per-step failure"
);
let after = sub.preflight("default", DB).await.unwrap();
assert!(
!after.iter().any(|a| a.id == "0003_atomic"),
"a rolled-back step must not be recorded"
);
let fix = sql_step("0003_atomic", "CREATE TABLE gadget (id int primary key)");
assert!(
matches!(apply(&sub, &fix, 2).await, SubstrateStepOutcome::Applied),
"retry after rollback applies cleanly"
);
assert!(matches!(
apply(&sub, &ext_step("0004_citext", "citext"), 3).await,
SubstrateStepOutcome::Applied
));
assert!(
matches!(
apply(&sub, &ext_step("0005_dblink", "dblink"), 4).await,
SubstrateStepOutcome::Failed(_)
),
"a non-allowlisted extension is a per-step failure"
);
assert!(matches!(
apply(
&sub,
&sql_step("0005_raw", "CREATE EXTENSION IF NOT EXISTS pgcrypto"),
4
)
.await,
SubstrateStepOutcome::Failed(_)
));
assert!(matches!(
apply(
&sub,
&sql_step("0006_txn", "BEGIN; CREATE TABLE sneaky (x int); COMMIT;"),
4
)
.await,
SubstrateStepOutcome::Failed(_)
));
let ddl = sub.owner_ddl("default", DB).await.unwrap();
assert!(matches!(
ddl.exec("SELECT * FROM boatramp_migrations.schema_migrations")
.await
.unwrap_err(),
MigrateDdlError::LedgerProtected
));
assert!(matches!(
ddl.exec("BEGIN; CREATE TABLE x(i int); COMMIT")
.await
.unwrap_err(),
MigrateDdlError::TxnControl
));
assert!(matches!(
ddl.exec("CREATE TABLE y(i int); COMMIT-- sneak")
.await
.unwrap_err(),
MigrateDdlError::TxnControl
));
assert!(matches!(
ddl.exec("DROP TABLE /*x*/ boatramp_migrations.schema_migrations")
.await
.unwrap_err(),
MigrateDdlError::LedgerProtected
));
ddl.exec("CREATE TABLE IF NOT EXISTS owner_made (n int)")
.await
.unwrap();
ddl.exec("INSERT INTO owner_made (n) VALUES (7)")
.await
.unwrap();
let rows = ddl
.query("SELECT n FROM owner_made ORDER BY n")
.await
.unwrap();
assert_eq!(rows.rows.len(), 1, "the owner sees the row it just wrote");
let _ = ddl.exec("DROP TABLE owner_made").await;
let baselined = sql_step("0007_baselined", "CREATE TABLE baselined (id int)");
sub.record(
"default",
DB,
&baselined,
4,
&baselined.content_hash(),
LedgerOrigin::Baseline,
)
.await
.unwrap();
let after = sub.preflight("default", DB).await.unwrap();
let row = after
.iter()
.find(|a| a.id == "0007_baselined")
.expect("baselined row present");
assert_eq!(row.origin, "baseline", "a baselined row is marked as such");
{
use boatramp_storage::sql_sqlx::{ExternalSqlKind, ExternalSqlOptions, connect};
let url = std::env::var("BOATRAMP_TEST_PG_URL").unwrap();
let c = connect(ExternalSqlKind::Postgres, &ExternalSqlOptions::new(url)).unwrap();
let exists = c
.run_query("SELECT to_regclass('public.baselined') IS NOT NULL AS e")
.await
.unwrap();
assert!(
matches!(
exists.rows.first().and_then(|r| r.first()),
Some(boatramp_core::sql::SqlValue::Boolean(false))
| Some(boatramp_core::sql::SqlValue::Null)
),
"baseline records without running the step"
);
}
{
use boatramp_storage::sql_sqlx::{ExternalSqlKind, ExternalSqlOptions, connect};
let url = std::env::var("BOATRAMP_TEST_PG_URL").unwrap();
let c = connect(ExternalSqlKind::Postgres, &ExternalSqlOptions::new(url)).unwrap();
for stmt in [
"DROP SCHEMA IF EXISTS boatramp_migrations CASCADE",
"DROP TABLE IF EXISTS widget",
"DROP TABLE IF EXISTS gadget",
"DROP TABLE IF EXISTS baselined",
] {
let _ = c.run_script(stmt).await;
}
}
println!("MIGRATE RUNNER LEDGER OK [postgres]");
}