use std::path::{Path, PathBuf};
use clap::Subcommand;
use serde::{Deserialize, Serialize};
use crate::client::ControlPlane;
#[derive(Debug, thiserror::Error)]
pub enum Error {
#[error(transparent)]
Client(#[from] crate::client::ClientError),
#[error("rendering the migration result failed: {0}")]
Render(#[source] serde_json::Error),
#[error("migration step `{id}` failed: {error}")]
StepFailed {
id: String,
error: String,
},
#[error("pass exactly one of --file <manifest.json> or --dir <migrations/>")]
BundleSource,
#[error("reading {path}: {source}")]
Read {
path: String,
#[source]
source: std::io::Error,
},
#[error("`{name}`: unrecognized migration file — expected .sql, .notx.sql, .ext, or .fn.json")]
UnknownStep {
name: String,
},
#[error("duplicate migration id `{id}` (files `{a}` and `{b}`)")]
DuplicateId {
id: String,
a: String,
b: String,
},
#[error("parsing function step `{name}`: {source}")]
FunctionStepParse {
name: String,
#[source]
source: serde_json::Error,
},
#[error("serializing the assembled bundle: {0}")]
Assemble(#[source] serde_json::Error),
#[error("{0}")]
InvalidDb(String),
}
fn validate_db(db: &str) -> Result<()> {
boatramp_core::project::validate_resource_name("database", db).map_err(|err| {
if db.is_empty() {
Error::InvalidDb(format!(
"{err}; {}",
boatramp_core::project::EMPTY_DB_NAME_CURE
))
} else {
Error::InvalidDb(format!(
"{err}; {}",
boatramp_core::project::INVALID_NAME_CURE
))
}
})
}
type Result<T> = std::result::Result<T, Error>;
#[derive(Debug, clap::Args)]
pub struct MigrateArgs {
#[command(subcommand)]
command: MigrateCommand,
}
#[derive(Debug, clap::Args)]
struct BundleSource {
#[arg(long, short = 'f', conflicts_with = "dir")]
file: Option<PathBuf>,
#[arg(long, short = 'd', conflicts_with = "file")]
dir: Option<PathBuf>,
}
#[derive(Debug, Subcommand)]
enum MigrateCommand {
Apply {
#[arg(long, default_value = boatramp_core::project::DEFAULT_DB_NAME)]
db: String,
#[command(flatten)]
source: BundleSource,
#[arg(long)]
json: bool,
},
DryRun {
#[arg(long, default_value = boatramp_core::project::DEFAULT_DB_NAME)]
db: String,
#[command(flatten)]
source: BundleSource,
#[arg(long)]
json: bool,
},
Baseline {
#[arg(long, default_value = boatramp_core::project::DEFAULT_DB_NAME)]
db: String,
#[command(flatten)]
source: BundleSource,
#[arg(long)]
up_to: Option<String>,
#[arg(long)]
json: bool,
},
Status {
#[arg(long, default_value = boatramp_core::project::DEFAULT_DB_NAME)]
db: String,
#[arg(long)]
json: bool,
},
}
pub async fn run(args: MigrateArgs, cp: &ControlPlane) -> Result<()> {
match args.command {
MigrateCommand::Apply { db, source, json } => {
validate_db(&db)?;
let bundle = resolve_bundle(cp, &source).await?;
let report = cp.migrate_trigger(&db, "apply", &bundle, None).await?;
render_report(&report, json)?;
fail_if_step_failed(report)?;
}
MigrateCommand::DryRun { db, source, json } => {
validate_db(&db)?;
let bundle = resolve_bundle(cp, &source).await?;
let report = cp.migrate_trigger(&db, "dry-run", &bundle, None).await?;
render_report(&report, json)?;
fail_if_step_failed(report)?;
}
MigrateCommand::Baseline {
db,
source,
up_to,
json,
} => {
validate_db(&db)?;
let bundle = resolve_bundle(cp, &source).await?;
let report = cp
.migrate_trigger(&db, "baseline", &bundle, up_to.as_deref())
.await?;
render_report(&report, json)?;
fail_if_step_failed(report)?;
}
MigrateCommand::Status { db, json } => {
validate_db(&db)?;
let status = cp.migrate_status(&db).await?;
render_status(&status, json)?;
}
}
Ok(())
}
#[derive(Debug, PartialEq)]
enum PickedSource<'a> {
File(&'a Path),
Dir(&'a Path),
}
fn pick_source(source: &BundleSource) -> Result<PickedSource<'_>> {
match (&source.file, &source.dir) {
(Some(file), None) => Ok(PickedSource::File(file)),
(None, Some(dir)) => Ok(PickedSource::Dir(dir)),
_ => Err(Error::BundleSource),
}
}
async fn resolve_bundle(cp: &ControlPlane, source: &BundleSource) -> Result<String> {
match pick_source(source)? {
PickedSource::File(file) => Ok(cp.put_file_blob(file).await?),
PickedSource::Dir(dir) => Ok(cp.put_bytes_blob(assemble_dir(dir)?).await?),
}
}
#[derive(Debug, Serialize, PartialEq)]
struct Step {
id: String,
#[serde(skip_serializing_if = "Option::is_none")]
sql: Option<String>,
#[serde(skip_serializing_if = "std::ops::Not::not")]
no_transaction: bool,
#[serde(skip_serializing_if = "Option::is_none")]
extension: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
function: Option<FunctionStep>,
}
#[derive(Debug, Serialize, Deserialize, PartialEq)]
struct FunctionStep {
name: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
version: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
args: Option<String>,
}
#[derive(Debug, Serialize, PartialEq)]
struct Bundle {
steps: Vec<Step>,
}
fn assemble_dir(dir: &Path) -> Result<Vec<u8>> {
let read = |p: &Path| -> Result<Vec<std::fs::DirEntry>> {
let mut entries: Vec<_> = std::fs::read_dir(p)
.map_err(|source| Error::Read {
path: p.display().to_string(),
source,
})?
.collect::<std::io::Result<Vec<_>>>()
.map_err(|source| Error::Read {
path: p.display().to_string(),
source,
})?;
entries.sort_by_key(std::fs::DirEntry::file_name);
Ok(entries)
};
let slurp = |path: &Path| -> Result<String> {
std::fs::read_to_string(path).map_err(|source| Error::Read {
path: path.display().to_string(),
source,
})
};
let mut steps: Vec<Step> = Vec::new();
let mut seen: std::collections::HashMap<String, String> = std::collections::HashMap::new();
for entry in read(dir)? {
let path = entry.path();
if path.is_dir() {
continue;
}
let name = entry.file_name().to_string_lossy().into_owned();
if name.starts_with('.') {
continue;
}
let (id, step) = classify_step(&name, &path, &slurp)?;
if let Some(prev) = seen.insert(id.clone(), name.clone()) {
return Err(Error::DuplicateId {
id,
a: prev,
b: name,
});
}
steps.push(step);
}
serde_json::to_vec(&Bundle { steps }).map_err(Error::Assemble)
}
fn classify_step(
name: &str,
path: &Path,
slurp: &dyn Fn(&Path) -> Result<String>,
) -> Result<(String, Step)> {
if let Some(id) = name.strip_suffix(".notx.sql") {
return Ok((
id.to_string(),
Step {
id: id.to_string(),
sql: Some(slurp(path)?),
no_transaction: true,
extension: None,
function: None,
},
));
}
if let Some(id) = name.strip_suffix(".sql") {
return Ok((
id.to_string(),
Step {
id: id.to_string(),
sql: Some(slurp(path)?),
no_transaction: false,
extension: None,
function: None,
},
));
}
if let Some(id) = name.strip_suffix(".ext") {
return Ok((
id.to_string(),
Step {
id: id.to_string(),
sql: None,
no_transaction: false,
extension: Some(slurp(path)?.trim().to_string()),
function: None,
},
));
}
if let Some(id) = name.strip_suffix(".fn.json") {
let body = slurp(path)?;
let function: FunctionStep =
serde_json::from_str(&body).map_err(|source| Error::FunctionStepParse {
name: name.to_string(),
source,
})?;
return Ok((
id.to_string(),
Step {
id: id.to_string(),
sql: None,
no_transaction: false,
extension: None,
function: Some(function),
},
));
}
Err(Error::UnknownStep {
name: name.to_string(),
})
}
fn render_report(report: &boatramp_core::sql::MigrationReport, json: bool) -> Result<()> {
if json {
println!(
"{}",
serde_json::to_string_pretty(report).map_err(Error::Render)?
);
return Ok(());
}
print_ids("applied", &report.newly_applied, &report.kinds);
print_ids("already applied", &report.already_applied, &report.kinds);
print_ids("pending", &report.pending, &report.kinds);
if let Some(failure) = &report.failed {
let kind = report
.kinds
.get(&failure.id)
.map(String::as_str)
.unwrap_or("?");
println!("FAILED at `{}` [{kind}]: {}", failure.id, failure.error);
}
if report.newly_applied.is_empty()
&& report.already_applied.is_empty()
&& report.pending.is_empty()
&& report.failed.is_none()
{
println!("nothing to do");
}
Ok(())
}
fn print_ids(label: &str, ids: &[String], kinds: &std::collections::BTreeMap<String, String>) {
if ids.is_empty() {
return;
}
println!("{label} ({}):", ids.len());
for id in ids {
let kind = kinds.get(id).map(String::as_str).unwrap_or("?");
println!(" {id} [{kind}]");
}
}
fn render_status(status: &boatramp_core::sql::MigrationStatus, json: bool) -> Result<()> {
if json {
println!(
"{}",
serde_json::to_string_pretty(status).map_err(Error::Render)?
);
return Ok(());
}
if status.applied.is_empty() {
println!("no migrations applied");
return Ok(());
}
println!(
"{:<4} {:<28} {:<10} {:<9} applied_at",
"#", "id", "kind", "origin"
);
for m in &status.applied {
println!(
"{:<4} {:<28} {:<10} {:<9} {}",
m.ordinal, m.id, m.kind, m.origin, m.applied_at
);
}
Ok(())
}
fn fail_if_step_failed(report: boatramp_core::sql::MigrationReport) -> Result<()> {
if let Some(failure) = report.failed {
return Err(Error::StepFailed {
id: failure.id,
error: failure.error,
});
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use boatramp_core::sql::{
AppliedMigration, MigrationFailure, MigrationReport, MigrationStatus,
};
use clap::Parser;
#[derive(Parser)]
struct Cli {
#[command(subcommand)]
cmd: Cmd,
}
#[derive(Subcommand)]
enum Cmd {
Migrate(MigrateArgs),
}
fn parse(argv: &[&str]) -> std::result::Result<MigrateCommand, clap::Error> {
let cli = Cli::try_parse_from(std::iter::once("boatramp").chain(argv.iter().copied()))?;
let Cmd::Migrate(args) = cli.cmd;
Ok(args.command)
}
#[test]
fn apply_and_baseline_flags_parse() {
match parse(&["migrate", "apply", "-f", "m.json"]) {
Ok(MigrateCommand::Apply { db, source, json }) => {
assert_eq!(db, boatramp_core::project::DEFAULT_DB_NAME);
assert_eq!(source.file, Some(PathBuf::from("m.json")));
assert_eq!(source.dir, None);
assert!(!json);
}
other => panic!("expected apply, got {other:?}"),
}
match parse(&["migrate", "apply", "-d", "migrations"]) {
Ok(MigrateCommand::Apply { source, .. }) => {
assert_eq!(source.dir, Some(PathBuf::from("migrations")));
assert_eq!(source.file, None);
}
other => panic!("expected apply -d, got {other:?}"),
}
assert!(matches!(
parse(&[
"migrate", "dry-run", "--db", "main", "-f", "m.json", "--json"
]),
Ok(MigrateCommand::DryRun { .. })
));
match parse(&[
"migrate",
"baseline",
"-f",
"m.json",
"--up-to",
"0003_seed",
]) {
Ok(MigrateCommand::Baseline { up_to, .. }) => {
assert_eq!(up_to.as_deref(), Some("0003_seed"));
}
other => panic!("expected baseline, got {other:?}"),
}
match parse(&["migrate", "status", "--db", "main"]) {
Ok(MigrateCommand::Status { db, json }) => {
assert_eq!(db, "main");
assert!(!json);
}
other => panic!("expected status, got {other:?}"),
}
assert!(parse(&["migrate", "apply", "-f", "m.json", "-d", "migrations"]).is_err());
}
#[test]
fn validate_db_gates_unsafe_names_client_side() {
assert!(validate_db(boatramp_core::project::DEFAULT_DB_NAME).is_ok());
assert!(validate_db("main").is_ok());
for bad in ["", "a/b", "..", "a b"] {
assert!(validate_db(bad).is_err(), "{bad:?} should be rejected");
}
}
#[test]
fn empty_db_name_rejection_carries_the_cure() {
let Err(Error::InvalidDb(msg)) = validate_db("") else {
panic!("empty --db must be rejected");
};
assert!(
msg.contains(boatramp_core::project::EMPTY_DB_NAME_CURE),
"empty-name error must carry the cure, got: {msg}"
);
let Err(Error::InvalidDb(other)) = validate_db("a/b") else {
panic!("`a/b` must be rejected");
};
assert!(
!other.contains(boatramp_core::project::EMPTY_DB_NAME_CURE),
"a non-empty invalid name must not carry the empty-name cure, got: {other}"
);
}
#[test]
fn bundle_source_requires_exactly_one() {
let neither = BundleSource {
file: None,
dir: None,
};
assert!(matches!(pick_source(&neither), Err(Error::BundleSource)));
let file = BundleSource {
file: Some(PathBuf::from("m.json")),
dir: None,
};
assert!(matches!(pick_source(&file), Ok(PickedSource::File(_))));
let dir = BundleSource {
file: None,
dir: Some(PathBuf::from("migrations")),
};
assert!(matches!(pick_source(&dir), Ok(PickedSource::Dir(_))));
}
#[test]
fn a_failed_step_is_a_non_zero_exit() {
let clean = MigrationReport {
newly_applied: vec!["0001_init".into()],
..MigrationReport::default()
};
assert!(fail_if_step_failed(clean).is_ok());
let failed = MigrationReport {
newly_applied: vec!["0001_init".into()],
failed: Some(MigrationFailure {
id: "0002_add_col".into(),
error: "relation already exists".into(),
}),
..MigrationReport::default()
};
match fail_if_step_failed(failed) {
Err(Error::StepFailed { id, error }) => {
assert_eq!(id, "0002_add_col");
assert!(error.contains("already exists"));
}
other => panic!("expected StepFailed, got {other:?}"),
}
}
#[test]
fn renderers_are_infallible_on_representative_payloads() {
let report = MigrationReport {
newly_applied: vec!["0001_init".into()],
already_applied: vec!["0000_base".into()],
pending: vec![],
failed: None,
kinds: [
("0001_init".to_string(), "sql".to_string()),
("0000_base".to_string(), "function".to_string()),
]
.into_iter()
.collect(),
};
render_report(&report, false).unwrap();
render_report(&report, true).unwrap();
render_report(&MigrationReport::default(), false).unwrap();
let status = MigrationStatus {
applied: vec![AppliedMigration {
id: "0001_init".into(),
ordinal: 0,
content_hash: "abc".into(),
kind: "sql".into(),
applied_at: "2026-09-22T00:00:00Z".into(),
origin: "apply".into(),
}],
};
render_status(&status, false).unwrap();
render_status(&status, true).unwrap();
render_status(&MigrationStatus::default(), false).unwrap();
}
fn name_as_body(p: &Path) -> Result<String> {
Ok(p.file_name().unwrap().to_string_lossy().into_owned())
}
#[test]
fn classify_step_maps_suffix_to_kind() {
let go = |name: &str| classify_step(name, Path::new(name), &name_as_body);
let (id, step) = go("0003_idx.notx.sql").unwrap();
assert_eq!(id, "0003_idx");
assert!(step.no_transaction);
assert_eq!(step.sql.as_deref(), Some("0003_idx.notx.sql"));
assert!(step.extension.is_none() && step.function.is_none());
let (id, step) = go("0001_init.sql").unwrap();
assert_eq!(id, "0001_init");
assert!(!step.no_transaction);
assert!(step.sql.is_some());
let (id, step) = go("0002_pgcrypto.ext").unwrap();
assert_eq!(id, "0002_pgcrypto");
assert_eq!(step.extension.as_deref(), Some("0002_pgcrypto.ext"));
assert!(matches!(go("README.md"), Err(Error::UnknownStep { .. })));
}
#[test]
fn classify_fn_step_parses_json_body() {
let body =
|_: &Path| Ok(r#"{"name":"backfill","version":"v3","args":"{\"n\":5}"}"#.to_string());
let (id, step) = classify_step("0004_backfill.fn.json", Path::new("x"), &body).unwrap();
assert_eq!(id, "0004_backfill");
assert_eq!(
step.function,
Some(FunctionStep {
name: "backfill".into(),
version: Some("v3".into()),
args: Some("{\"n\":5}".into()),
})
);
let bad = |_: &Path| Ok("not json".to_string());
assert!(matches!(
classify_step("0005_x.fn.json", Path::new("x"), &bad),
Err(Error::FunctionStepParse { .. })
));
}
#[test]
fn assemble_dir_orders_by_filename_and_rejects_dupe_ids() {
let root = std::env::temp_dir().join(format!("br-migrate-asm-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&root);
std::fs::create_dir_all(&root).unwrap();
std::fs::write(root.join("0002_pgcrypto.ext"), "pgcrypto\n").unwrap();
std::fs::write(root.join("0001_init.sql"), "CREATE TABLE t (id int);").unwrap();
std::fs::write(root.join("0003_backfill.fn.json"), r#"{"name":"backfill"}"#).unwrap();
std::fs::write(root.join(".gitkeep"), "").unwrap();
let bytes = assemble_dir(&root).unwrap();
let bundle: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
let steps = bundle["steps"].as_array().unwrap();
let ids: Vec<&str> = steps.iter().map(|s| s["id"].as_str().unwrap()).collect();
assert_eq!(ids, ["0001_init", "0002_pgcrypto", "0003_backfill"]);
assert!(steps[0].get("sql").is_some());
assert_eq!(steps[1]["extension"], "pgcrypto");
assert_eq!(steps[2]["function"]["name"], "backfill");
assert!(steps[0].get("no_transaction").is_none());
std::fs::write(root.join("0001_init.ext"), "dup").unwrap();
assert!(matches!(
assemble_dir(&root),
Err(Error::DuplicateId { id, .. }) if id == "0001_init"
));
let _ = std::fs::remove_dir_all(&root);
}
}