use futures_util::TryStreamExt;
use sqlx::{ConnectOptions, Connection, Executor, pool::PoolOptions};
use sqlx_turso_driver::{Turso, TursoConnectOptions, TursoPool, connect_pool};
use std::time::{Duration, Instant};
use tokio::time::timeout;
type TestResult = Result<(), Box<dyn std::error::Error>>;
const BOUND: Duration = Duration::from_secs(2);
#[tokio::test]
async fn file_pool_shares_committed_data() -> TestResult {
let directory = tempfile::tempdir()?;
let options = TursoConnectOptions::file(directory.path().join("pool 中文.db"))?;
let pool: TursoPool = timeout(BOUND, connect_pool(options.clone(), 2)).await??;
let mut first = timeout(BOUND, pool.acquire()).await??;
first.execute("CREATE TABLE items(value INTEGER)").await?;
first
.execute("CREATE TEMP TABLE local_state(value INTEGER)")
.await?;
let mut second = timeout(BOUND, pool.acquire()).await??;
assert!(
second
.execute("INSERT INTO local_state VALUES (1)")
.await
.is_err(),
"distinct slots must not clone the same physical connection"
);
assert_eq!(pool.size(), 2);
first.execute("INSERT INTO items VALUES (7)").await?; assert_eq!(
sqlx::query_scalar::<Turso, i64>("SELECT value FROM items")
.fetch_one(&mut *second)
.await?,
7
);
let held = Instant::now();
assert!(
timeout(Duration::from_millis(30), pool.acquire())
.await
.is_err()
);
assert!(held.elapsed() < BOUND);
assert_eq!(pool.size(), 2);
drop(first);
drop(second);
timeout(BOUND, pool.close()).await?;
assert!(matches!(pool.acquire().await, Err(sqlx::Error::PoolClosed)));
let mut reopened = options.connect().await?;
assert_eq!(
sqlx::query_scalar::<Turso, i64>("SELECT value FROM items")
.fetch_one(&mut reopened)
.await?,
7
);
reopened.close().await?;
Ok(())
}
#[tokio::test]
async fn memory_pool_rejects_multiple_connections() -> TestResult {
for capacity in [0, 2] {
assert!(matches!(
connect_pool(TursoConnectOptions::memory(), capacity).await,
Err(sqlx::Error::Configuration(_))
));
}
let directory = tempfile::tempdir()?;
assert!(matches!(
connect_pool(
TursoConnectOptions::file(directory.path().join("zero.db"))?,
0
)
.await,
Err(sqlx::Error::Configuration(_))
));
let first = connect_pool(TursoConnectOptions::memory(), 1).await?;
let second = connect_pool(TursoConnectOptions::memory(), 1).await?;
first.execute("CREATE TABLE items(value INTEGER)").await?;
first.execute("INSERT INTO items VALUES (7)").await?;
assert_eq!(
sqlx::query_scalar::<Turso, i64>("SELECT value FROM items")
.fetch_one(&first)
.await?,
7
);
assert!(
sqlx::query_scalar::<Turso, i64>("SELECT value FROM items")
.fetch_one(&second)
.await
.is_err()
);
timeout(BOUND, first.close()).await?;
timeout(BOUND, second.close()).await?;
Ok(())
}
#[tokio::test]
async fn dropped_stream_releases_pool_slot() -> TestResult {
for memory in [true, false] {
let directory = tempfile::tempdir()?;
let options = if memory {
TursoConnectOptions::memory()
} else {
TursoConnectOptions::file(directory.path().join("stream.db"))?
};
let pool = connect_pool(options, 1).await?;
pool.execute("CREATE TABLE items(value INTEGER)").await?;
pool.execute("INSERT INTO items VALUES (7), (8), (9)")
.await?;
let mut stream =
sqlx::query_scalar::<Turso, i64>("SELECT value FROM items ORDER BY value").fetch(&pool);
assert_eq!(timeout(BOUND, stream.try_next()).await??, Some(7));
assert!(
timeout(Duration::from_millis(30), pool.acquire())
.await
.is_err()
);
drop(stream);
let mut connection = timeout(BOUND, pool.acquire()).await??;
connection.ping().await?;
assert_eq!(
connection
.execute("INSERT INTO items VALUES (10)")
.await?
.rows_affected(),
1
);
assert_eq!(
sqlx::query_scalar::<Turso, i64>("SELECT count(*) FROM items")
.fetch_one(&mut *connection)
.await?,
4
);
drop(connection);
timeout(BOUND, pool.close()).await?;
assert!(matches!(pool.acquire().await, Err(sqlx::Error::PoolClosed)));
}
Ok(())
}
#[tokio::test]
async fn acquire_timeout_does_not_leak() -> TestResult {
let directory = tempfile::tempdir()?;
let pool = PoolOptions::<Turso>::new()
.max_connections(1)
.acquire_timeout(Duration::from_millis(40))
.connect_with(TursoConnectOptions::file(
directory.path().join("timeout.db"),
)?)
.await?;
let mut held = pool.acquire().await?;
held.execute("CREATE TABLE items(value INTEGER)").await?;
held.execute("CREATE TEMP TABLE local_state(value INTEGER)")
.await?;
held.execute("INSERT INTO local_state VALUES (11)").await?;
assert!(matches!(
timeout(BOUND, pool.acquire()).await?,
Err(sqlx::Error::PoolTimedOut)
));
assert!(
timeout(Duration::from_millis(10), pool.acquire())
.await
.is_err()
);
drop(held);
let mut next = timeout(BOUND, pool.acquire()).await??;
assert_eq!(
sqlx::query_scalar::<Turso, i64>("SELECT value FROM local_state")
.fetch_one(&mut *next)
.await?,
11
);
next.execute("INSERT INTO items VALUES (7)").await?;
drop(next);
assert_eq!(
timeout(
BOUND,
sqlx::query_scalar::<Turso, i64>("SELECT value FROM items").fetch_one(&pool)
)
.await??,
7
);
timeout(BOUND, pool.close()).await?;
assert!(matches!(pool.acquire().await, Err(sqlx::Error::PoolClosed)));
Ok(())
}
#[tokio::test]
async fn memory_connection_loss_is_explicit() -> TestResult {
let pool = connect_pool(TursoConnectOptions::memory(), 1).await?;
pool.execute("CREATE TABLE items(value INTEGER)").await?;
let connection = pool.acquire().await?;
connection.close().await?;
let error = timeout(BOUND, pool.acquire()).await?.unwrap_err();
assert!(matches!(error, sqlx::Error::Configuration(_)), "{error}");
assert!(error.to_string().contains("memory"), "{error}");
timeout(BOUND, pool.close()).await?;
Ok(())
}
#[tokio::test]
async fn ping_cleans_abandoned_reader() -> TestResult {
let mut connection = TursoConnectOptions::memory().connect().await?;
connection.ping().await?;
connection
.execute("CREATE TABLE items(value INTEGER)")
.await?;
connection
.execute("INSERT INTO items VALUES (7), (8)")
.await?;
let mut stream =
sqlx::query_scalar::<Turso, i64>("SELECT value FROM items").fetch(&mut connection);
assert_eq!(stream.try_next().await?, Some(7));
drop(stream);
connection.ping().await?;
connection.execute("INSERT INTO items VALUES (9)").await?;
connection.close().await?;
Ok(())
}
#[tokio::test]
async fn successful_flush_preserves_clean_pool_connection() -> TestResult {
for memory in [false, true] {
let directory = tempfile::tempdir()?;
let options = if memory {
TursoConnectOptions::memory()
} else {
TursoConnectOptions::file(directory.path().join("flush.db"))?
};
let pool = connect_pool(options, 1).await?;
let mut connection = timeout(BOUND, pool.acquire()).await??;
connection
.execute("CREATE TABLE items(value INTEGER)")
.await?;
connection
.execute("INSERT INTO items VALUES (7), (8)")
.await?;
connection
.execute("CREATE TEMP TABLE local_state(value INTEGER)")
.await?;
connection
.execute("INSERT INTO local_state VALUES (11)")
.await?;
let mut stream = sqlx::query_scalar::<Turso, i64>("SELECT value FROM items ORDER BY value")
.fetch(&mut *connection);
assert_eq!(timeout(BOUND, stream.try_next()).await??, Some(7));
drop(stream);
assert!(connection.should_flush());
connection.flush().await?;
assert!(!connection.should_flush());
connection.ping().await?;
drop(connection);
let mut next = timeout(BOUND, pool.acquire()).await??;
assert_eq!(
sqlx::query_scalar::<Turso, i64>("SELECT value FROM local_state")
.fetch_one(&mut *next)
.await?,
11
);
assert_eq!(
sqlx::query_scalar::<Turso, i64>("SELECT sum(value) FROM items")
.fetch_one(&mut *next)
.await?,
15
);
next.flush().await?;
drop(next);
timeout(BOUND, pool.close()).await?;
}
Ok(())
}
#[tokio::test]
async fn busy_timeout_reaches_engine() -> TestResult {
for (duration, milliseconds) in [(Duration::from_millis(37), 37), (Duration::ZERO, 0)] {
let mut connection = TursoConnectOptions::memory()
.busy_timeout(duration)
.connect()
.await?;
assert_eq!(
sqlx::query_scalar::<Turso, i64>("PRAGMA busy_timeout")
.fetch_one(&mut connection)
.await?,
milliseconds
);
connection.close().await?;
}
Ok(())
}
#[tokio::test]
async fn busy_writer_fails_finitely_then_pool_recovers() -> TestResult {
let directory = tempfile::tempdir()?;
let path = directory.path().join("busy.db");
let pool = connect_pool(
TursoConnectOptions::file(&path)?.busy_timeout(Duration::from_millis(40)),
1,
)
.await?;
pool.execute("CREATE TABLE items(value INTEGER)").await?;
let database = turso::Builder::new_local(path.to_str().unwrap())
.build()
.await?;
let writer = database.connect()?;
writer.execute("BEGIN IMMEDIATE", ()).await?;
writer.execute("INSERT INTO items VALUES (1)", ()).await?;
let started = Instant::now();
let error = timeout(BOUND, pool.execute("INSERT INTO items VALUES (2)"))
.await?
.unwrap_err();
assert!(matches!(error, sqlx::Error::Database(_)), "{error}");
assert!(started.elapsed() < BOUND);
writer.execute("ROLLBACK", ()).await?;
assert_eq!(
timeout(BOUND, pool.execute("INSERT INTO items VALUES (3)"))
.await??
.rows_affected(),
1
);
assert_eq!(
sqlx::query_scalar::<Turso, i64>("SELECT sum(value) FROM items")
.fetch_one(&pool)
.await?,
3
);
timeout(BOUND, pool.close()).await?;
Ok(())
}
#[tokio::test]
async fn cancelled_busy_query_releases_clean_connection() -> TestResult {
let directory = tempfile::tempdir()?;
let path = directory.path().join("cancel.db");
let pool = connect_pool(
TursoConnectOptions::file(&path)?.busy_timeout(Duration::from_millis(200)),
1,
)
.await?;
let mut connection = pool.acquire().await?;
connection
.execute("CREATE TABLE items(value INTEGER)")
.await?;
connection
.execute("CREATE TEMP TABLE local_state(value INTEGER)")
.await?;
connection
.execute("INSERT INTO local_state VALUES (11)")
.await?;
let database = turso::Builder::new_local(path.to_str().unwrap())
.build()
.await?;
let writer = database.connect()?;
writer.execute("BEGIN IMMEDIATE", ()).await?;
writer.execute("INSERT INTO items VALUES (1)", ()).await?;
assert!(
timeout(
Duration::from_millis(10),
connection.execute("INSERT INTO items VALUES (2)")
)
.await
.is_err()
);
drop(connection);
let mut next = timeout(BOUND, pool.acquire()).await??;
assert_eq!(
sqlx::query_scalar::<Turso, i64>("SELECT value FROM local_state")
.fetch_one(&mut *next)
.await?,
11
);
writer.execute("ROLLBACK", ()).await?;
assert_eq!(
next.execute("INSERT INTO items VALUES (3)")
.await?
.rows_affected(),
1
);
assert_eq!(
sqlx::query_scalar::<Turso, i64>("SELECT sum(value) FROM items")
.fetch_one(&mut *next)
.await?,
3
);
drop(next);
timeout(BOUND, pool.close()).await?;
Ok(())
}
#[tokio::test]
async fn open_raw_transaction_is_not_returned_as_clean() -> TestResult {
let directory = tempfile::tempdir()?;
let pool = connect_pool(
TursoConnectOptions::file(directory.path().join("unclean.db"))?,
1,
)
.await?;
let mut connection = pool.acquire().await?;
connection
.execute("CREATE TABLE items(value INTEGER)")
.await?;
connection.begin().await?.rollback().await?;
connection.execute("BEGIN IMMEDIATE").await?;
connection.execute("INSERT INTO items VALUES (1)").await?;
assert!(connection.ping().await.is_err());
drop(connection);
let mut next = timeout(BOUND, pool.acquire()).await??;
assert_eq!(
sqlx::query_scalar::<Turso, i64>("SELECT count(*) FROM items")
.fetch_one(&mut *next)
.await?,
0
);
next.execute("INSERT INTO items VALUES (7)").await?;
drop(next);
timeout(BOUND, pool.close()).await?;
Ok(())
}
#[tokio::test]
async fn shutdown_rejects_waiting_acquire() -> TestResult {
let pool = connect_pool(TursoConnectOptions::memory(), 1).await?;
let held = pool.acquire().await?;
let mut close = Box::pin(pool.close());
assert!(
timeout(Duration::from_millis(20), &mut close)
.await
.is_err()
);
assert!(pool.is_closed());
assert!(matches!(
timeout(BOUND, pool.acquire()).await?,
Err(sqlx::Error::PoolClosed)
));
drop(held);
timeout(BOUND, close).await?;
assert_eq!(pool.size(), 0);
Ok(())
}