sqlx-turso-driver 0.0.1

An asynchronous SQLx driver for embedded Turso databases
use futures_util::TryStreamExt;
use sqlx::{Column, ConnectOptions, Executor, FromRow, Row, SqlSafeStr, Statement, TypeInfo};
use sqlx_turso_driver::{Turso, TursoConnectOptions, TursoRow, TursoTypeInfo};

type TestResult = Result<(), Box<dyn std::error::Error>>;

// Catches rejection, fabricated expression types and execution during preparation.
#[tokio::test]
async fn prepared_statement_can_execute() -> TestResult {
    let mut connection = TursoConnectOptions::memory().connect().await?;
    connection
        .execute("CREATE TABLE items(value INTEGER, label TEXT, data BLOB, price REAL)")
        .await?;
    let statement = connection
        .prepare("SELECT value, label, data, price, ?1 AS expression FROM items".into_sql_str())
        .await?;
    assert_eq!(statement.columns().len(), 5);
    assert_eq!(statement.column("value").ordinal(), 0);
    assert_eq!(statement.column("label").type_info(), &TursoTypeInfo::Text);
    assert_eq!(statement.column(0).type_info(), &TursoTypeInfo::Integer);
    assert_eq!(statement.column(2).type_info(), &TursoTypeInfo::Blob);
    assert_eq!(statement.column(3).type_info(), &TursoTypeInfo::Real);
    // The public Turso wrapper exposes neither parameter types nor counts.
    assert!(statement.parameters().is_none());
    assert_eq!(statement.column(4).type_info().name(), "UNKNOWN");
    assert!(matches!(
        statement.try_column("missing"),
        Err(sqlx::Error::ColumnNotFound(_))
    ));
    assert!(matches!(
        statement.try_column(5),
        Err(sqlx::Error::ColumnIndexOutOfBounds { .. })
    ));
    let insert = connection
        .prepare("INSERT INTO items(value) VALUES (?1)".into_sql_str())
        .await?;
    assert_eq!(
        sqlx::query_scalar::<Turso, i64>("SELECT count(*) FROM items")
            .fetch_one(&mut connection)
            .await?,
        0
    );
    insert.query().bind(7_i64).execute(&mut connection).await?;
    let row = statement
        .query()
        .bind("中文")
        .fetch_one(&mut connection)
        .await?;
    assert_eq!(row.try_get::<i64, _>("value")?, 7);
    assert_eq!(row.try_get::<String, _>("expression")?, "中文");
    assert!(
        connection
            .prepare_with("SELECT ?1".into_sql_str(), &[TursoTypeInfo::Text])
            .await
            .is_err()
    );
    Ok(())
}

// Catches column-presence heuristics dropping actual RETURNING write counts.
#[tokio::test]
async fn returning_reports_actual_affected_rows() -> TestResult {
    let mut connection = TursoConnectOptions::memory().connect().await?;
    connection
        .execute("CREATE TABLE items(value INTEGER)")
        .await?;
    let result = sqlx::query::<Turso>("INSERT INTO items VALUES (1), (2) RETURNING value")
        .execute(&mut connection)
        .await?;
    assert_eq!(result.rows_affected(), 2);
    let result = sqlx::query::<Turso>("UPDATE items SET value = value + 1 RETURNING value")
        .execute(&mut connection)
        .await?;
    assert_eq!(result.rows_affected(), 2);
    let result = sqlx::query::<Turso>("DELETE FROM items WHERE value = 2 RETURNING value")
        .execute(&mut connection)
        .await?;
    assert_eq!(result.rows_affected(), 1);
    Ok(())
}

// Catches stale counts after reads/DDL, zero-match writes and incorrect conflict counts.
#[tokio::test]
async fn affected_rows_match_completed_modifications() -> TestResult {
    let mut connection = TursoConnectOptions::memory().connect().await?;
    for (sql, expected) in [
        ("CREATE TABLE items(value INTEGER UNIQUE)", 0),
        ("INSERT INTO items VALUES (1), (2), (3)", 3),
        ("SELECT * FROM items", 0),
        ("SELECT 1 WHERE 0", 0),
        ("CREATE TABLE other(value INTEGER)", 0),
        ("UPDATE items SET value = value + 10 WHERE value < 3", 2),
        ("UPDATE items SET value = value WHERE value = 11", 1),
        ("UPDATE items SET value = 99 WHERE value = 999", 0),
        ("INSERT OR IGNORE INTO items VALUES (11), (13)", 1),
        ("DELETE FROM items WHERE value > 10", 3),
        ("DELETE FROM items WHERE value > 10", 0),
        // Pinned Turso counts the removed sqlite_schema row for DROP TABLE.
        ("DROP TABLE other", 1),
    ] {
        assert_eq!(
            sqlx::query::<Turso>(sql)
                .execute(&mut connection)
                .await?
                .rows_affected(),
            expected,
            "{sql}"
        );
    }
    // Independent engine characterization: n_change includes schema deletion,
    // unlike connection-level changes(), and must not be silently normalized.
    let database = turso::Builder::new_local(":memory:").build().await?;
    let raw = database.connect()?;
    raw.execute("CREATE TABLE other(value INTEGER)", ()).await?;
    let mut drop_statement = raw.prepare("DROP TABLE other").await?;
    drop_statement.execute(()).await?;
    assert_eq!(drop_statement.n_change(), 1);
    Ok(())
}

#[derive(Debug, PartialEq)]
struct Record {
    id: i64,
    label: String,
}
impl FromRow<'_, TursoRow> for Record {
    fn from_row(row: &TursoRow) -> Result<Self, sqlx::Error> {
        Ok(Self {
            id: row.try_get("id")?,
            label: row.try_get("label")?,
        })
    }
}

// Catches broken FromRow/fetch modes, name indexing and rows borrowing dead engine buffers.
#[tokio::test]
async fn typed_fetch_modes_own_their_rows() -> TestResult {
    let mut connection = TursoConnectOptions::memory().connect().await?;
    connection
        .execute("CREATE TABLE items(id INTEGER, label TEXT)")
        .await?;
    sqlx::query::<Turso>("INSERT INTO items VALUES (?, ?), (?, ?)")
        .bind(1_i64)
        .bind("甲")
        .bind(2_i64)
        .bind("乙")
        .execute(&mut connection)
        .await?;
    let records = sqlx::query_as::<Turso, Record>("SELECT * FROM items ORDER BY id")
        .fetch_all(&mut connection)
        .await?;
    assert_eq!(
        records,
        [
            Record {
                id: 1,
                label: "甲".into()
            },
            Record {
                id: 2,
                label: "乙".into()
            }
        ]
    );
    let record = sqlx::query_as::<Turso, Record>("SELECT * FROM items WHERE id = ?")
        .bind(2_i64)
        .fetch_one(&mut connection)
        .await?;
    assert_eq!(
        record,
        Record {
            id: 2,
            label: "乙".into()
        }
    );
    assert!(
        sqlx::query_as::<Turso, Record>("SELECT * FROM items WHERE 0")
            .fetch_optional(&mut connection)
            .await?
            .is_none()
    );
    assert!(matches!(
        sqlx::query::<Turso>("SELECT * FROM items WHERE 0")
            .fetch_one(&mut connection)
            .await,
        Err(sqlx::Error::RowNotFound)
    ));
    let mut stream = sqlx::query::<Turso>("SELECT * FROM items ORDER BY id").fetch(&mut connection);
    let first = stream.try_next().await?.unwrap();
    let second = stream.try_next().await?.unwrap();
    assert!(stream.try_next().await?.is_none());
    drop(stream);
    connection
        .execute("UPDATE items SET label = 'changed'")
        .await?;
    drop(connection);
    assert_eq!(first.try_get::<&str, _>("label")?, "甲");
    assert_eq!(second.try_get::<&str, _>(1)?, "乙");
    assert!(matches!(
        first.try_get::<i64, _>("missing"),
        Err(sqlx::Error::ColumnNotFound(_))
    ));
    assert!(matches!(
        first.try_get::<i64, _>(2),
        Err(sqlx::Error::ColumnIndexOutOfBounds { index: 2, len: 2 })
    ));
    Ok(())
}

// Catches fabricating UNIQUE subtypes from text and hiding failed writes as successes.
#[tokio::test]
async fn constraint_failure_preserves_source_without_guessing_subtype() -> TestResult {
    let mut connection = TursoConnectOptions::memory().connect().await?;
    connection
        .execute("CREATE TABLE items(id INTEGER UNIQUE)")
        .await?;
    sqlx::query::<Turso>("INSERT INTO items VALUES (?)")
        .bind(7_i64)
        .execute(&mut connection)
        .await?;
    let error = sqlx::query::<Turso>("INSERT INTO items VALUES (?)")
        .bind(7_i64)
        .execute(&mut connection)
        .await
        .unwrap_err();
    let sqlx::Error::Database(error) = error else {
        panic!("expected database error")
    };
    assert_eq!(error.kind(), sqlx::error::ErrorKind::Other);
    assert!(error.code().is_none());
    assert!(error.constraint().is_none());
    let source = std::error::Error::source(&*error)
        .unwrap()
        .downcast_ref::<turso::Error>()
        .unwrap();
    assert!(matches!(source, turso::Error::Constraint(_)));
    assert_eq!(error.message(), source.to_string());
    assert_eq!(
        sqlx::query_scalar::<Turso, i64>("SELECT count(*) FROM items")
            .fetch_one(&mut connection)
            .await?,
        1
    );
    Ok(())
}

// Catches explicit preparation bypassing the D1 real-parser batch guard.
#[tokio::test]
async fn prepared_batches_are_rejected_without_side_effects() -> TestResult {
    let mut connection = TursoConnectOptions::memory().connect().await?;
    assert!(
        connection
            .prepare("CREATE TABLE ignored(x); SELECT 1".into_sql_str())
            .await
            .is_err()
    );
    assert!(
        sqlx::query::<Turso>("SELECT * FROM ignored")
            .fetch_all(&mut connection)
            .await
            .is_err()
    );
    let statement = connection
        .prepare("SELECT ?1 AS \"中文;列\"; -- ; trailing comment\n".into_sql_str())
        .await?;
    assert_eq!(
        statement
            .query()
            .bind(7_i64)
            .fetch_one(&mut connection)
            .await?
            .try_get::<i64, _>("中文;列")?,
        7
    );
    Ok(())
}