use rusqlite::{Connection, OptionalExtension};
use sha2::{Digest, Sha256};
use thiserror::Error;
#[derive(Error, Debug)]
#[non_exhaustive]
pub enum MigrationError {
#[error("Migration failed: {0}")]
Failed(String),
#[error("Database error: {0}")]
Database(#[from] rusqlite::Error),
}
pub fn run_migrations(conn: &mut Connection) -> Result<(), MigrationError> {
conn.execute(
"CREATE TABLE IF NOT EXISTS _migrations (
id INTEGER PRIMARY KEY,
name TEXT NOT NULL UNIQUE,
checksum TEXT,
applied_at TEXT NOT NULL DEFAULT (datetime('now'))
)",
[],
)?;
ensure_migration_checksum_column(conn)?;
let migrations = get_migrations();
for (name, sql) in migrations {
if name == "027_vector_search" && !cfg!(feature = "vector") {
continue;
}
if name == "028_bm25_search" && !fts5_available(conn) {
continue;
}
let checksum = compute_migration_checksum(sql);
let existing_checksum: Option<String> = conn
.query_row("SELECT checksum FROM _migrations WHERE name = ?", [name], |row| row.get(0))
.optional()?;
if let Some(existing) = existing_checksum {
if !existing.is_empty() && existing != checksum {
return Err(MigrationError::Failed(format!(
"migration checksum mismatch for {name}: expected {checksum}, found {existing}",
)));
}
if existing.is_empty() {
conn.execute(
"UPDATE _migrations SET checksum = ? WHERE name = ?",
rusqlite::params![checksum, name],
)?;
}
} else {
let tx = conn.transaction()?;
tx.execute_batch(sql)?;
tx.execute(
"INSERT INTO _migrations (name, checksum) VALUES (?, ?)",
rusqlite::params![name, checksum],
)?;
tx.commit()?;
}
}
Ok(())
}
fn fts5_available(conn: &Connection) -> bool {
conn.query_row("SELECT sqlite_compileoption_used('ENABLE_FTS5')", [], |row| {
row.get::<_, i32>(0)
})
.optional()
.ok()
.flatten()
.unwrap_or(0)
== 1
}
fn ensure_migration_checksum_column(conn: &Connection) -> Result<(), MigrationError> {
let has_checksum = conn
.prepare("PRAGMA table_info(_migrations)")?
.query_map([], |row| row.get::<_, String>(1))?
.collect::<std::result::Result<Vec<_>, _>>()?
.into_iter()
.any(|column| column == "checksum");
if !has_checksum {
conn.execute("ALTER TABLE _migrations ADD COLUMN checksum TEXT", [])?;
}
Ok(())
}
fn compute_migration_checksum(sql: &str) -> String {
let digest = Sha256::digest(sql.as_bytes());
digest.iter().map(|b| format!("{b:02x}")).collect()
}
#[must_use]
pub fn known_migration_names() -> Vec<&'static str> {
get_migrations().into_iter().map(|(name, _)| name).collect()
}
#[must_use]
pub fn latest_known_migration() -> &'static str {
get_migrations().last().map_or("", |(name, _)| *name)
}
fn get_migrations() -> Vec<(&'static str, &'static str)> {
vec![
("001_initial_schema", include_str!("../migrations/001_initial_schema.sql")),
("002_inventory", include_str!("../migrations/002_inventory.sql")),
("003_returns", include_str!("../migrations/003_returns.sql")),
("004_manufacturing", include_str!("../migrations/004_manufacturing.sql")),
("005_shipments", include_str!("../migrations/005_shipments.sql")),
(
"006_payments_warranties_po_invoices",
include_str!("../migrations/006_payments_warranties_po_invoices.sql"),
),
("007_carts", include_str!("../migrations/007_carts.sql")),
("008_multi_currency", include_str!("../migrations/008_multi_currency.sql")),
("009_tax", include_str!("../migrations/009_tax.sql")),
("010_promotions", include_str!("../migrations/010_promotions.sql")),
("011_subscriptions", include_str!("../migrations/011_subscriptions.sql")),
("012_versioning", include_str!("../migrations/012_versioning.sql")),
("013_quality", include_str!("../migrations/013_quality.sql")),
("014_lots", include_str!("../migrations/014_lots.sql")),
("015_serials", include_str!("../migrations/015_serials.sql")),
("016_warehouse", include_str!("../migrations/016_warehouse.sql")),
("017_receiving", include_str!("../migrations/017_receiving.sql")),
("018_fulfillment", include_str!("../migrations/018_fulfillment.sql")),
("019_accounts_payable", include_str!("../migrations/019_accounts_payable.sql")),
("020_cost_accounting", include_str!("../migrations/020_cost_accounting.sql")),
("021_credit", include_str!("../migrations/021_credit.sql")),
("022_backorder", include_str!("../migrations/022_backorder.sql")),
("023_accounts_receivable", include_str!("../migrations/023_accounts_receivable.sql")),
("024_general_ledger", include_str!("../migrations/024_general_ledger.sql")),
("025_performance_indexes", include_str!("../migrations/025_performance_indexes.sql")),
("026_idempotency_keys", include_str!("../migrations/026_idempotency_keys.sql")),
("027_vector_search", include_str!("../migrations/027_vector_search.sql")),
("028_bm25_search", include_str!("../migrations/028_bm25_search.sql")),
("029_x402_a2a", include_str!("../migrations/029_x402_a2a.sql")),
("030_x402_credits", include_str!("../migrations/030_x402_credits.sql")),
("031_erc8004", include_str!("../migrations/031_erc8004.sql")),
("032_custom_objects", include_str!("../migrations/032_custom_objects.sql")),
("033_orders_cart_id", include_str!("../migrations/033_orders_cart_id.sql")),
("034_x402_nonce_integrity", include_str!("../migrations/034_x402_nonce_integrity.sql")),
("035_x402_pqc", include_str!("../migrations/035_x402_pqc.sql")),
(
"036_fix_updated_at_triggers",
include_str!("../migrations/036_fix_updated_at_triggers.sql"),
),
("037_commerce_entities", include_str!("../migrations/037_commerce_entities.sql")),
(
"038_fix_location_inventory_trigger",
include_str!("../migrations/038_fix_location_inventory_trigger.sql"),
),
("039_b2b_erp_entities", include_str!("../migrations/039_b2b_erp_entities.sql")),
("040_supplier_skus", include_str!("../migrations/040_supplier_skus.sql")),
("041_vendor_returns", include_str!("../migrations/041_vendor_returns.sql")),
("042_vendor_credits", include_str!("../migrations/042_vendor_credits.sql")),
("043_payment_obligations", include_str!("../migrations/043_payment_obligations.sql")),
("044_price_levels", include_str!("../migrations/044_price_levels.sql")),
("045_prepayments", include_str!("../migrations/045_prepayments.sql")),
("046_price_schedules", include_str!("../migrations/046_price_schedules.sql")),
("047_activity_logs", include_str!("../migrations/047_activity_logs.sql")),
("048_integration_mappings", include_str!("../migrations/048_integration_mappings.sql")),
("049_inbound_shipments", include_str!("../migrations/049_inbound_shipments.sql")),
("050_purgatory", include_str!("../migrations/050_purgatory.sql")),
("051_print_stations", include_str!("../migrations/051_print_stations.sql")),
("052_edi_documents", include_str!("../migrations/052_edi_documents.sql")),
(
"053_integration_field_mappings",
include_str!("../migrations/053_integration_field_mappings.sql"),
),
("054_topology_snapshots", include_str!("../migrations/054_topology_snapshots.sql")),
("055_stock_snapshots", include_str!("../migrations/055_stock_snapshots.sql")),
("056_rewards", include_str!("../migrations/056_rewards.sql")),
("057_loyalty_tiers", include_str!("../migrations/057_loyalty_tiers.sql")),
("058_wishlist_item_fields", include_str!("../migrations/058_wishlist_item_fields.sql")),
("059_fraud", include_str!("../migrations/059_fraud.sql")),
("060_fixed_assets_revrec", include_str!("../migrations/060_fixed_assets_revrec.sql")),
("061_gl_auto_posting_flags", include_str!("../migrations/061_gl_auto_posting_flags.sql")),
("062_cycle_counts", include_str!("../migrations/062_cycle_counts.sql")),
("063_gl_fx_revaluation", include_str!("../migrations/063_gl_fx_revaluation.sql")),
("064_http_idempotency", include_str!("../migrations/064_http_idempotency.sql")),
("065_zone_shipping_methods", include_str!("../migrations/065_zone_shipping_methods.sql")),
("066_search_configs", include_str!("../migrations/066_search_configs.sql")),
]
}