use fathomdb_schema::{
check_migration_accretion, migrate, migrate_with_event_sink, migrate_with_steps, Migration,
MigrationAccretionError, MigrationError, SCHEMA_VERSION,
};
use rusqlite::Connection;
use std::process::Command;
use std::sync::Once;
fn register_sqlite_vec_once() {
static REGISTER: Once = Once::new();
REGISTER.call_once(|| unsafe {
let entrypoint: unsafe extern "C" fn(
*mut rusqlite::ffi::sqlite3,
*mut *const std::os::raw::c_char,
*const rusqlite::ffi::sqlite3_api_routines,
) -> std::os::raw::c_int = std::mem::transmute(sqlite_vec::sqlite3_vec_init as *const ());
rusqlite::ffi::sqlite3_auto_extension(Some(entrypoint));
});
}
fn user_version(conn: &Connection) -> u32 {
conn.query_row("PRAGMA user_version", [], |row| row.get::<_, u32>(0)).unwrap()
}
fn set_user_version(conn: &Connection, version: u32) {
conn.pragma_update(None, "user_version", version).unwrap();
}
#[test]
fn ac_046a_applies_ordered_migrations_to_current_version() {
register_sqlite_vec_once();
let conn = Connection::open_in_memory().unwrap();
set_user_version(&conn, 1);
let report = migrate(&conn).unwrap();
assert_eq!(report.schema_version_before, 1);
assert_eq!(report.schema_version_after, SCHEMA_VERSION);
assert_eq!(user_version(&conn), SCHEMA_VERSION);
assert_eq!(report.migration_steps.len(), 9);
assert!(report.migration_steps.iter().all(|step| !step.failed));
}
#[test]
fn ac_046b_success_report_contains_step_ids_and_durations() {
register_sqlite_vec_once();
let conn = Connection::open_in_memory().unwrap();
set_user_version(&conn, 1);
let report = migrate(&conn).unwrap();
let step_ids: Vec<u32> = report.migration_steps.iter().map(|step| step.step_id).collect();
assert_eq!(step_ids, vec![2, 3, 4, 5, 6, 7, 8, 9, 10]);
assert!(report.migration_steps.iter().all(|step| step.duration_ms.is_some()));
}
#[test]
fn ac_046b_success_emits_structured_step_events() {
register_sqlite_vec_once();
let conn = Connection::open_in_memory().unwrap();
set_user_version(&conn, 1);
let mut events = Vec::new();
migrate_with_event_sink(&conn, fathomdb_schema::MIGRATIONS, |event| {
events.push(event.clone());
})
.unwrap();
let step_ids: Vec<u32> = events.iter().map(|step| step.step_id).collect();
assert_eq!(step_ids, vec![2, 3, 4, 5, 6, 7, 8, 9, 10]);
assert!(events.iter().all(|step| step.duration_ms.is_some()));
assert!(events.iter().all(|step| !step.failed));
}
#[test]
fn ac_046c_and_ac_070_failed_step_reports_failure_and_preserves_user_version() {
let conn = Connection::open_in_memory().unwrap();
set_user_version(&conn, 1);
let migrations = [Migration {
step_id: 2,
sql: "CREATE TABLE _poison(id INTEGER PRIMARY KEY); SELECT * FROM missing_table",
}];
let err = migrate_with_steps(&conn, &migrations).expect_err("poison migration must fail");
match err {
MigrationError::MigrationError(report) => {
assert_eq!(report.schema_version_before, 1);
assert_eq!(report.schema_version_current, 1);
assert_eq!(report.migration_steps.len(), 1);
assert_eq!(report.migration_steps[0].step_id, 2);
assert!(report.migration_steps[0].failed);
assert!(report.migration_steps[0].duration_ms.is_some());
}
other => panic!("expected MigrationError, got {other:?}"),
}
assert_eq!(user_version(&conn), 1);
}
#[test]
fn ac_046c_failed_migration_emits_failed_step_event() {
let conn = Connection::open_in_memory().unwrap();
set_user_version(&conn, 1);
let migrations = [Migration {
step_id: 2,
sql: "CREATE TABLE _poison(id INTEGER PRIMARY KEY); SELECT * FROM missing_table",
}];
let mut events = Vec::new();
let err = migrate_with_event_sink(&conn, &migrations, |event| events.push(event.clone()))
.expect_err("poison migration must fail");
assert!(matches!(err, MigrationError::MigrationError(_)));
assert_eq!(events.len(), 1);
assert_eq!(events[0].step_id, 2);
assert!(events[0].failed);
assert!(events[0].duration_ms.is_some());
}
#[test]
fn ac_049_accretion_guard_accepts_current_migrations_and_rejects_violator() {
check_migration_accretion(
"003_add_profile.sql",
"CREATE TABLE x(id INTEGER); -- MIGRATION-ACCRETION-EXEMPTION: bootstrap profile table",
)
.unwrap();
let err = check_migration_accretion("004_bad.sql", "ALTER TABLE x ADD COLUMN y TEXT")
.expect_err("adding a column without removal or exemption must fail");
assert_eq!(err, MigrationAccretionError { offender: "004_bad.sql".to_string() });
}
#[test]
fn phase9_pack_b_migration_008_adds_source_id_columns_and_indexes() {
register_sqlite_vec_once();
let conn = Connection::open_in_memory().unwrap();
set_user_version(&conn, 1);
migrate(&conn).unwrap();
let nodes_has_source_id: bool = conn
.prepare("PRAGMA table_info(canonical_nodes)")
.unwrap()
.query_map([], |row| row.get::<_, String>(1))
.unwrap()
.filter_map(Result::ok)
.any(|name| name == "source_id");
assert!(nodes_has_source_id, "canonical_nodes.source_id must be present after migration 8");
let edges_has_source_id: bool = conn
.prepare("PRAGMA table_info(canonical_edges)")
.unwrap()
.query_map([], |row| row.get::<_, String>(1))
.unwrap()
.filter_map(Result::ok)
.any(|name| name == "source_id");
assert!(edges_has_source_id, "canonical_edges.source_id must be present after migration 8");
let nodes_idx: u64 = conn
.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type='index' AND name=?1",
["canonical_nodes_source_id_idx"],
|row| row.get(0),
)
.unwrap();
assert_eq!(nodes_idx, 1);
let edges_idx: u64 = conn
.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type='index' AND name=?1",
["canonical_edges_source_id_idx"],
|row| row.get(0),
)
.unwrap();
assert_eq!(edges_idx, 1);
assert_eq!(user_version(&conn), SCHEMA_VERSION);
}
#[test]
fn ac_049_repo_linter_accepts_actual_migrations_and_names_violator() {
let repo =
std::path::Path::new(env!("CARGO_MANIFEST_DIR")).ancestors().nth(4).expect("repo root");
let script = repo.join("scripts/agent-lint-migrations.sh");
let ok = Command::new(&script).current_dir(repo).output().expect("run linter");
assert!(
ok.status.success(),
"stdout={} stderr={}",
String::from_utf8_lossy(&ok.stdout),
String::from_utf8_lossy(&ok.stderr)
);
let fixture = repo
.join("src/rust/crates/fathomdb-schema/tests/fixtures/migrations/accretion_violator.sql");
let failed = Command::new(&script)
.arg(&fixture)
.current_dir(repo)
.output()
.expect("run linter on fixture");
assert!(!failed.status.success());
assert!(String::from_utf8_lossy(&failed.stderr).contains("accretion_violator.sql"));
}