fathomdb-schema 0.7.1

FathomDB schema — versioned migration registry and bootstrap (leaf crate).
Documentation
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;

// Pack 1 (step 9) creates a vec0 virtual table; register sqlite-vec
// once per test-binary so the migration step can execute.
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"));
}