use sqlx::migrate::{Migrate, MigrateDatabase, MigrateError, Migration, MigrationType, Migrator};
use sqlx::{ConnectOptions, Connection, Executor, SqlSafeStr};
use sqlx_turso_driver::{Turso, TursoConnectOptions, TursoConnection, connect_pool};
use std::{future::Future, time::Duration};
#[path = "support/migration_process.rs"]
mod migration_process;
type TestResult = Result<(), Box<dyn std::error::Error>>;
const TABLE: &str = "_sqlx_migrations";
async fn bounded<T>(future: impl Future<Output = T>) -> T {
tokio::time::timeout(Duration::from_secs(10), future)
.await
.expect("migration deadline")
}
fn migration(version: i64, sql: &'static str) -> Migration {
Migration::new(
version,
"D5 fixture".into(),
MigrationType::Simple,
sql.into_sql_str(),
false,
)
}
fn migrator(sql: &'static str) -> Migrator {
Migrator::with_migrations(vec![migration(1, sql)])
}
async fn schema_count(conn: &mut TursoConnection, name: &str) -> Result<i64, sqlx::Error> {
sqlx::query_scalar("SELECT count(*) FROM sqlite_schema WHERE name = ?")
.bind(name)
.fetch_one(conn)
.await
}
#[tokio::test]
async fn direct_sdk_ddl_rollback_is_isolated() -> TestResult {
bounded(async {
let directory = tempfile::tempdir()?;
let path = directory.path().join("sdk-gate.db");
let db = turso::Builder::new_local(path.to_str().unwrap())
.build()
.await?;
let owner = db.connect()?;
let observer_db = turso::Builder::new_local(path.to_str().unwrap())
.build()
.await?;
let observer = observer_db.connect()?;
owner.execute("BEGIN", ()).await?;
owner
.execute("CREATE TABLE partial(id INTEGER PRIMARY KEY)", ())
.await?;
owner.execute("INSERT INTO partial VALUES (1)", ()).await?;
let mut own = owner.query("SELECT count(*) FROM partial", ()).await?;
assert_eq!(
own.next().await?.unwrap().get_value(0)?,
turso::Value::Integer(1)
);
drop(own);
let mut hidden = observer
.query(
"SELECT count(*) FROM sqlite_schema WHERE name='partial'",
(),
)
.await?;
assert_eq!(
hidden.next().await?.unwrap().get_value(0)?,
turso::Value::Integer(0)
);
drop(hidden);
let error = owner
.execute("INSERT INTO partial VALUES (1)", ())
.await
.unwrap_err();
assert!(
error.to_string().to_lowercase().contains("unique"),
"{error}"
);
owner.execute("ROLLBACK", ()).await?;
let mut absent = owner
.query(
"SELECT count(*) FROM sqlite_schema WHERE name='partial'",
(),
)
.await?;
assert_eq!(
absent.next().await?.unwrap().get_value(0)?,
turso::Value::Integer(0)
);
drop(absent);
drop(owner);
drop(observer);
drop(db);
drop(observer_db);
let mut reopened = TursoConnectOptions::file(path)?.connect().await?;
assert_eq!(schema_count(&mut reopened, "partial").await?, 0);
reopened.close().await?;
TestResult::Ok(())
})
.await
}
#[tokio::test]
async fn applied_once_checksum_is_stable() -> TestResult {
bounded(async {
for memory in [false, true] {
let directory = tempfile::tempdir()?;
let options = if memory { TursoConnectOptions::memory() } else {
TursoConnectOptions::file(directory.path().join("迁移 data.db"))?
};
let pool = connect_pool(options.clone(), 1).await?;
let migrations = Migrator::with_migrations(vec![
migration(1, "-- 注释 ;\nCREATE TABLE items(id INTEGER PRIMARY KEY, value TEXT); /* ; */ INSERT INTO items VALUES (1, 'hello;世界') RETURNING id;"),
migration(2, "CREATE INDEX items_value ON items(value); INSERT INTO items VALUES (2, 'second'); -- final ;\n"),
]);
migrations.run(&pool).await?;
migrations.run(&pool).await?;
let mut conn = pool.acquire().await?;
let applied = conn.list_applied_migrations(TABLE).await?;
assert_eq!(applied.iter().map(|m| m.version).collect::<Vec<_>>(), vec![1, 2]);
assert_eq!(applied[0].checksum, migrations.migrations[0].checksum);
assert_eq!(applied[1].checksum, migrations.migrations[1].checksum);
assert_eq!(sqlx::query_scalar::<Turso, i64>("SELECT count(*) FROM items").fetch_one(&mut *conn).await?, 2);
assert_eq!(sqlx::query_scalar::<Turso, String>("SELECT value FROM items WHERE id=1").fetch_one(&mut *conn).await?, "hello;世界");
assert_eq!(conn.dirty_version(TABLE).await?, None);
drop(conn);
pool.close().await;
if !memory {
let mut reopened = options.connect().await?;
migrations.run(&mut reopened).await?;
assert_eq!(reopened.list_applied_migrations(TABLE).await?.len(), 2);
reopened.close().await?;
}
}
TestResult::Ok(())
}).await
}
#[tokio::test]
async fn changed_checksum_is_rejected() -> TestResult {
bounded(async {
let directory = tempfile::tempdir()?;
let options = TursoConnectOptions::file(directory.path().join("checksum.db"))?;
let pool = connect_pool(options.clone(), 1).await?;
migrator("CREATE TABLE original(id INTEGER);")
.run(&pool)
.await?;
assert!(matches!(
migrator("CREATE TABLE changed(id INTEGER);")
.run(&pool)
.await,
Err(MigrateError::VersionMismatch(1))
));
let mut next = pool.acquire().await?;
assert_eq!(schema_count(&mut next, "changed").await?, 0);
let mut competitor = options.connect().await?;
competitor.lock().await?;
competitor.unlock().await?;
competitor.close().await?;
drop(next);
pool.close().await;
TestResult::Ok(())
})
.await
}
#[tokio::test]
async fn failed_ddl_is_not_applied() -> TestResult {
bounded(async {
let directory = tempfile::tempdir()?;
let options = TursoConnectOptions::file(directory.path().join("failed.db"))?;
let pool = connect_pool(options.clone(), 1).await?;
let bad = migrator("CREATE TABLE partial(id INTEGER PRIMARY KEY); INSERT INTO partial VALUES (1); INSERT INTO partial VALUES (1); CREATE TABLE never_run(id INTEGER);");
let error = bad.run(&pool).await.unwrap_err();
assert!(matches!(&error, MigrateError::ExecuteMigration(sqlx::Error::Database(_), 1)), "{error:?}");
assert!(error.to_string().to_lowercase().contains("unique"), "DDL must execute before runtime failure: {error}");
let mut conn = pool.acquire().await?;
assert_eq!(schema_count(&mut conn, "partial").await?, 0);
assert_eq!(schema_count(&mut conn, "never_run").await?, 0);
assert!(conn.list_applied_migrations(TABLE).await?.is_empty());
assert_eq!(conn.dirty_version(TABLE).await?, None);
drop(conn);
pool.close().await;
let mut reopened = options.connect().await?;
assert_eq!(schema_count(&mut reopened, "partial").await?, 0);
assert!(reopened.list_applied_migrations(TABLE).await?.is_empty());
migrator("CREATE TABLE partial(id INTEGER PRIMARY KEY); INSERT INTO partial VALUES (2);").run(&mut reopened).await?;
assert_eq!(reopened.list_applied_migrations(TABLE).await?.len(), 1);
reopened.close().await?;
TestResult::Ok(())
}).await
}
#[tokio::test]
async fn metadata_failure_rolls_back_whole_version() -> TestResult {
bounded(async {
let directory = tempfile::tempdir()?;
let options = TursoConnectOptions::file(directory.path().join("metadata.db"))?;
let pool = connect_pool(options.clone(), 1).await?;
pool.execute("CREATE TABLE _sqlx_migrations(version BIGINT PRIMARY KEY CHECK(version < 2), description TEXT NOT NULL, installed_on TEXT DEFAULT CURRENT_TIMESTAMP, success BOOLEAN NOT NULL, checksum BLOB NOT NULL, execution_time BIGINT NOT NULL)").await?;
let m = Migrator::with_migrations(vec![migration(2, "CREATE TABLE metadata_partial(id INTEGER);")]);
assert!(matches!(m.run(&pool).await, Err(MigrateError::ExecuteMigration(sqlx::Error::Database(_), 2))));
let mut conn = pool.acquire().await?;
assert_eq!(schema_count(&mut conn, "metadata_partial").await?, 0);
assert!(conn.list_applied_migrations(TABLE).await?.is_empty());
drop(conn);
pool.close().await;
let mut reopened = options.connect().await?;
assert_eq!(schema_count(&mut reopened, "metadata_partial").await?, 0);
reopened.close().await?;
TestResult::Ok(())
}).await
}
#[tokio::test]
async fn dirty_version_is_rejected_without_recovery() -> TestResult {
bounded(async {
let directory = tempfile::tempdir()?;
let options = TursoConnectOptions::file(directory.path().join("dirty.db"))?;
let pool = connect_pool(options.clone(), 1).await?;
let mut conn = pool.acquire().await?;
conn.ensure_migrations_table(TABLE).await?;
conn.execute("INSERT INTO _sqlx_migrations(version, description, success, checksum, execution_time) VALUES (1, 'dirty', 0, x'00', 0)").await?;
assert_eq!(conn.dirty_version(TABLE).await?, Some(1));
drop(conn);
assert!(matches!(migrator("CREATE TABLE skipped(id INTEGER);").run(&pool).await, Err(MigrateError::Dirty(1))));
let mut next = pool.acquire().await?;
assert_eq!(schema_count(&mut next, "skipped").await?, 0);
assert_eq!(next.dirty_version(TABLE).await?, Some(1));
assert!(next.list_applied_migrations(TABLE).await?.is_empty());
let mut other = options.connect().await?;
other.lock().await?;
other.unlock().await?;
other.close().await?;
drop(next);
pool.close().await;
TestResult::Ok(())
}).await
}
#[tokio::test]
async fn unsupported_operations_return_errors() -> TestResult {
bounded(async {
let mut conn = TursoConnectOptions::memory().connect().await?;
conn.ensure_migrations_table(TABLE).await?;
let m = migration(1, "CREATE TABLE forbidden(id INTEGER);");
assert!(conn.revert(TABLE, &m).await.is_err());
assert!(matches!(
conn.skip(TABLE, &m).await,
Err(MigrateError::SkipNotSupported())
));
assert!(matches!(
conn.create_schema_if_not_exists("extra").await,
Err(MigrateError::CreateSchemasNotSupported(_))
));
assert!(conn.ensure_migrations_table("main.custom").await.is_err());
assert!(
conn.ensure_migrations_table("x; DROP TABLE _sqlx_migrations")
.await
.is_err()
);
conn.lock().await?;
let mut no_tx = m.clone();
no_tx.no_tx = true;
assert!(conn.apply(TABLE, &no_tx).await.is_err());
assert_eq!(schema_count(&mut conn, "forbidden").await?, 0);
let down = Migration::new(
1,
"down".into(),
MigrationType::ReversibleDown,
"DROP TABLE forbidden".into_sql_str(),
false,
);
assert!(conn.apply(TABLE, &down).await.is_err());
conn.unlock().await?;
conn.close().await?;
TestResult::Ok(())
})
.await
}
#[tokio::test]
async fn migration_cannot_escape_owned_transaction() -> TestResult {
bounded(async {
for sql in [
"CREATE TABLE escape(id INTEGER); COMMIT; INSERT INTO escape VALUES (1);",
"CREATE TABLE escape(id INTEGER); ROLLBACK;",
"CREATE TABLE escape(id INTEGER); PRAGMA journal_mode=OFF;",
"CREATE TABLE escape(id INTEGER); ATTACH ':memory:' AS extra;",
"CREATE TABLE escape(id INTEGER); SAVEPOINT nested;",
"CREATE TABLE escape(id INTEGER); VACUUM;",
] {
let mut conn = TursoConnectOptions::memory().connect().await?;
assert!(
migrator(sql).run(&mut conn).await.is_err(),
"accepted: {sql}"
);
assert_eq!(schema_count(&mut conn, "escape").await?, 0);
assert!(conn.list_applied_migrations(TABLE).await?.is_empty());
conn.close_hard().await?;
}
TestResult::Ok(())
})
.await
}
#[tokio::test]
async fn other_migrator_early_errors_discard_lock_owner() -> TestResult {
bounded(async {
let directory = tempfile::tempdir()?;
let options = TursoConnectOptions::file(directory.path().join("early.db"))?;
let pool = connect_pool(options.clone(), 1).await?;
migrator("CREATE TABLE original(id INTEGER);")
.run(&pool)
.await?;
assert!(matches!(
Migrator::with_migrations(vec![]).run(&pool).await,
Err(MigrateError::VersionMissing(1))
));
let conn = pool.acquire().await?;
drop(conn);
let mut schemas = migrator("CREATE TABLE original(id INTEGER);");
schemas.create_schema("extra");
assert!(matches!(
schemas.run(&pool).await,
Err(MigrateError::CreateSchemasNotSupported(_))
));
let conn = pool.acquire().await?;
let mut other = options.connect().await?;
other.lock().await?;
other.unlock().await?;
other.close().await?;
drop(conn);
pool.close().await;
TestResult::Ok(())
})
.await
}
#[tokio::test]
async fn canonical_alias_lock_and_owner_lifecycle() -> TestResult {
bounded(async {
let directory = tempfile::tempdir()?;
let path = directory.path().join("actual data.db");
let options = TursoConnectOptions::file(&path)?;
let mut owner = options.connect().await?;
#[cfg(unix)]
let alias = {
let alias = directory.path().join("链接 alias.db");
std::os::unix::fs::symlink(&path, &alias)?;
alias
};
#[cfg(not(unix))]
let alias = path.clone();
let mut other = TursoConnectOptions::file(alias)?.connect().await?;
owner.lock().await?;
assert!(other.lock().await.is_err());
owner.close().await?;
other.lock().await?;
other.unlock().await?;
let mut owner = options.connect().await?;
owner.lock().await?;
owner.close_hard().await?;
other.lock().await?;
other.unlock().await?;
other.close().await?;
TestResult::Ok(())
})
.await
}
#[tokio::test]
async fn local_database_lifecycle_is_explicit() -> TestResult {
bounded(async {
let directory = tempfile::tempdir()?;
let path = directory.path().join("生命周期 db.db");
let url = url::Url::from_file_path(&path).unwrap().to_string();
assert!(!Turso::database_exists(&url).await?);
Turso::create_database(&url).await?;
assert!(Turso::database_exists(&url).await?);
let mut conn = TursoConnectOptions::file(&path)?.connect().await?;
conn.execute("CREATE TABLE retained(id INTEGER)").await?;
conn.close().await?;
Turso::create_database(&url).await?;
let mut reopened = TursoConnectOptions::file(path)?.connect().await?;
assert_eq!(schema_count(&mut reopened, "retained").await?, 1);
assert!(Turso::drop_database(&url).await.is_err());
assert!(Turso::force_drop_database(&url).await.is_err());
assert!(Turso::create_database("https://remote/db").await.is_err());
assert!(Turso::database_exists("turso-memory:").await.is_err());
reopened.close().await?;
TestResult::Ok(())
})
.await
}
#[tokio::test]
async fn apply_requires_lock_and_lock_misuse_errors() -> TestResult {
bounded(async {
let mut conn = TursoConnectOptions::memory().connect().await?;
conn.ensure_migrations_table(TABLE).await?;
let m = migration(1, "CREATE TABLE unlocked(id INTEGER);");
assert!(conn.apply(TABLE, &m).await.is_err());
let mut unlocked = Migrator::with_migrations(vec![m]);
unlocked.set_locking(false);
assert!(unlocked.run(&mut conn).await.is_err());
assert_eq!(schema_count(&mut conn, "unlocked").await?, 0);
conn.lock().await?;
assert!(conn.lock().await.is_err());
conn.unlock().await?;
assert!(conn.unlock().await.is_err());
conn.close().await?;
TestResult::Ok(())
})
.await
}
#[tokio::test]
async fn migration_lock_excludes_second_process() -> TestResult {
bounded(migration_process::exclusion_and_abnormal_exit()).await
}