Skip to main content

clt_database/turso/src/
lib.rs

1//! # Turso bindings for Rust
2//!
3//! Turso is an in-process SQL database engine, compatible with SQLite.
4//!
5//! ## Getting Started
6//!
7//! To get started, you first need to create a [`Database`] object and then open a [`Connection`] to it, which you use to query:
8//!
9//! ```rust,no_run
10//! # async fn run() {
11//! use turso::Builder;
12//!
13//! let db = Builder::new_local(":memory:").build().await.unwrap();
14//! let conn = db.connect().unwrap();
15//! conn.execute("CREATE TABLE IF NOT EXISTS users (email TEXT)", ()).await.unwrap();
16//! conn.execute("INSERT INTO users (email) VALUES ('alice@example.org')", ()).await.unwrap();
17//! # }
18//! ```
19//!
20//! You can also prepare statements with the [`Connection`] object and then execute the [`Statement`] objects:
21//!
22//! ```rust,no_run
23//! # async fn run() {
24//! # use turso::Builder;
25//! # let db = Builder::new_local(":memory:").build().await.unwrap();
26//! # let conn = db.connect().unwrap();
27//! let mut stmt = conn.prepare("SELECT * FROM users WHERE email = ?1").await.unwrap();
28//! let mut rows = stmt.query(["foo@example.com"]).await.unwrap();
29//! let row = rows.next().await.unwrap().unwrap();
30//! let value = row.get_value(0).unwrap();
31//! println!("Row: {:?}", value);
32//! # }
33//! ```
34
35#[cfg(all(clt_turso_feature = "mimalloc", not(target_family = "wasm"), not(miri)))]
36#[global_allocator]
37static GLOBAL: mimalloc::MiMalloc = mimalloc::MiMalloc;
38
39pub mod connection;
40pub mod params;
41mod rows;
42pub mod transaction;
43pub mod value;
44
45use crate::turso_sdk_kit::rsapi::TursoError;
46pub use connection::Connection;
47pub use value::Value;
48
49pub use params::params_from_iter;
50pub use params::IntoParams;
51
52use std::fmt::Debug;
53use std::future::Future;
54use std::sync::Arc;
55use std::sync::Mutex;
56use std::task::Poll;
57
58// Re-exports rows
59pub use crate::turso::rows::{Row, Rows};
60
61// Re-export turso_core
62pub use turso_core as core;
63
64/// Assert that a type implements both Send and Sync at compile time.
65/// Usage: assert_send_sync!(MyType);
66/// Usage: assert_send_sync!(Type1, Type2, Type3);
67macro_rules! assert_send_sync {
68    ($($t:ty),+ $(,)?) => {
69        #[cfg(clt_turso_tests)]
70        $(const _: () = {
71            const fn _assert_send<T: ?Sized + Send>() {}
72            const fn _assert_sync<T: ?Sized + Sync>() {}
73            _assert_send::<$t>();
74            _assert_sync::<$t>();
75        };)+
76    };
77}
78
79pub(crate) use assert_send_sync;
80
81#[derive(Debug, thiserror::Error)]
82pub enum Error {
83    #[error("SQL conversion failure: `{0}`")]
84    ToSqlConversionFailure(BoxError),
85    #[error("Query returned no rows")]
86    QueryReturnedNoRows,
87    #[error("Conversion failure: `{0}`")]
88    ConversionFailure(String),
89    #[error("{0}")]
90    Busy(String),
91    #[error("{0}")]
92    BusySnapshot(String),
93    #[error("{0}")]
94    Interrupt(String),
95    #[error("{0}")]
96    Error(String),
97    #[error("{0}")]
98    Misuse(String),
99    #[error("{0}")]
100    Constraint(String),
101    #[error("{0}")]
102    Readonly(String),
103    #[error("{0}")]
104    DatabaseFull(String),
105    #[error("{0}")]
106    NotAdb(String),
107    #[error("{0}")]
108    Corrupt(String),
109    #[error("I/O error ({1}): {0}")]
110    IoError(std::io::ErrorKind, &'static str),
111}
112
113impl From<crate::turso_sdk_kit::rsapi::TursoError> for Error {
114    fn from(value: crate::turso_sdk_kit::rsapi::TursoError) -> Self {
115        match value {
116            crate::turso_sdk_kit::rsapi::TursoError::Busy(err) => Error::Busy(err),
117            crate::turso_sdk_kit::rsapi::TursoError::BusySnapshot(err) => Error::BusySnapshot(err),
118            crate::turso_sdk_kit::rsapi::TursoError::Interrupt(err) => Error::Interrupt(err),
119            crate::turso_sdk_kit::rsapi::TursoError::Error(err) => Error::Error(err),
120            crate::turso_sdk_kit::rsapi::TursoError::Misuse(err) => Error::Misuse(err),
121            crate::turso_sdk_kit::rsapi::TursoError::Constraint(err) => Error::Constraint(err),
122            crate::turso_sdk_kit::rsapi::TursoError::Readonly(err) => Error::Readonly(err),
123            crate::turso_sdk_kit::rsapi::TursoError::DatabaseFull(err) => Error::DatabaseFull(err),
124            crate::turso_sdk_kit::rsapi::TursoError::NotAdb(err) => Error::NotAdb(err),
125            crate::turso_sdk_kit::rsapi::TursoError::Corrupt(err) => Error::Corrupt(err),
126            crate::turso_sdk_kit::rsapi::TursoError::IoError(kind, op) => Error::IoError(kind, op),
127        }
128    }
129}
130
131pub(crate) type BoxError = Box<dyn std::error::Error + Send + Sync>;
132
133pub type Result<T> = std::result::Result<T, Error>;
134pub type EncryptionOpts = crate::turso_sdk_kit::rsapi::EncryptionOpts;
135
136/// A builder for `Database`.
137pub struct Builder {
138    path: String,
139    enable_encryption: bool,
140    enable_attach: bool,
141    enable_custom_types: bool,
142    enable_index_method: bool,
143    enable_materialized_views: bool,
144    enable_vacuum: bool,
145    enable_generated_columns: bool,
146    enable_multiprocess_wal: bool,
147    enable_without_rowid: bool,
148    enable_mvcc_passive_checkpoint: bool,
149    vfs: Option<String>,
150    encryption_opts: Option<crate::turso_sdk_kit::rsapi::EncryptionOpts>,
151    io: Option<Arc<dyn turso_core::IO>>,
152}
153
154impl Builder {
155    /// Create a new local database.
156    pub fn new_local(path: &str) -> Self {
157        Self {
158            path: path.to_string(),
159            enable_encryption: false,
160            enable_attach: false,
161            enable_custom_types: false,
162            enable_index_method: false,
163            enable_materialized_views: false,
164            enable_vacuum: false,
165            enable_generated_columns: false,
166            enable_multiprocess_wal: false,
167            enable_without_rowid: false,
168            enable_mvcc_passive_checkpoint: false,
169            vfs: None,
170            encryption_opts: None,
171            io: None,
172        }
173    }
174
175    pub fn experimental_encryption(mut self, encryption_enabled: bool) -> Self {
176        self.enable_encryption = encryption_enabled;
177        self
178    }
179
180    pub fn with_encryption(mut self, opts: crate::turso_sdk_kit::rsapi::EncryptionOpts) -> Self {
181        self.encryption_opts = Some(opts);
182        self
183    }
184
185    /// Kept for backwards compatibility. Triggers are now always enabled.
186    pub fn experimental_triggers(self, _triggers_enabled: bool) -> Self {
187        self
188    }
189
190    pub fn experimental_attach(mut self, attach_enabled: bool) -> Self {
191        self.enable_attach = attach_enabled;
192        self
193    }
194
195    /// Kept for backwards compatibility. Strict tables are now always enabled.
196    pub fn experimental_strict(self, _strict_enabled: bool) -> Self {
197        self
198    }
199
200    pub fn experimental_custom_types(mut self, custom_types_enabled: bool) -> Self {
201        self.enable_custom_types = custom_types_enabled;
202        self
203    }
204
205    pub fn experimental_generated_columns(mut self, gencols_enabled: bool) -> Self {
206        self.enable_generated_columns = gencols_enabled;
207        self
208    }
209
210    pub fn experimental_index_method(mut self, index_method_enabled: bool) -> Self {
211        self.enable_index_method = index_method_enabled;
212        self
213    }
214
215    pub fn experimental_materialized_views(mut self, enabled: bool) -> Self {
216        self.enable_materialized_views = enabled;
217        self
218    }
219
220    pub fn experimental_vacuum(mut self, enabled: bool) -> Self {
221        self.enable_vacuum = enabled;
222        self
223    }
224
225    pub fn experimental_multiprocess_wal(mut self, enabled: bool) -> Self {
226        self.enable_multiprocess_wal = enabled;
227        self
228    }
229
230    pub fn experimental_without_rowid(mut self, enabled: bool) -> Self {
231        self.enable_without_rowid = enabled;
232        self
233    }
234
235    pub fn experimental_mvcc_passive_checkpoint(mut self, enabled: bool) -> Self {
236        self.enable_mvcc_passive_checkpoint = enabled;
237        self
238    }
239
240    pub fn with_io(mut self, vfs: String) -> Self {
241        self.vfs = Some(vfs);
242        self
243    }
244
245    /// Can pass custom IO implementation
246    pub fn with_io_impl(mut self, io: Arc<dyn turso_core::IO>) -> Self {
247        self.io = Some(io);
248        self
249    }
250
251    fn build_features_string(&self) -> Option<String> {
252        let mut features = Vec::new();
253        if self.enable_encryption {
254            features.push("encryption");
255        }
256        if self.enable_attach {
257            features.push("attach");
258        }
259        if self.enable_custom_types {
260            features.push("custom_types");
261        }
262        if self.enable_index_method {
263            features.push("index_method");
264        }
265        if self.enable_materialized_views {
266            features.push("views");
267        }
268        if self.enable_vacuum {
269            features.push("vacuum");
270        }
271        if self.enable_generated_columns {
272            features.push("generated_columns");
273        }
274        if self.enable_multiprocess_wal {
275            features.push("multiprocess_wal");
276        }
277        if self.enable_without_rowid {
278            features.push("without_rowid");
279        }
280        if self.enable_mvcc_passive_checkpoint {
281            features.push("mvcc_passive_checkpoint");
282        }
283        if features.is_empty() {
284            return None;
285        }
286        Some(features.join(","))
287    }
288
289    /// Build the database.
290    #[allow(unused_variables, clippy::arc_with_non_send_sync)]
291    pub async fn build(self) -> Result<Database> {
292        let features = self.build_features_string();
293        let db = crate::turso_sdk_kit::rsapi::TursoDatabase::new(
294            crate::turso_sdk_kit::rsapi::TursoDatabaseConfig {
295                path: self.path,
296                experimental_features: features,
297                async_io: true,
298                encryption: self.encryption_opts,
299                vfs: self.vfs,
300                io: self.io,
301                db_file: None,
302            },
303        );
304        while let Some(io_c) = db.open()?.io() {
305            // At this point IO must already be created
306            let io = db
307                .io()
308                .expect("IO must have been set on the first call to db open");
309            io_c.wait_async(io.as_ref())
310                .await
311                .map_err(TursoError::from)?;
312        }
313        Ok(Database { inner: db })
314    }
315}
316
317/// A database.
318///
319/// The `Database` object points to a database and allows you to connect to it
320#[derive(Clone)]
321pub struct Database {
322    inner: Arc<crate::turso_sdk_kit::rsapi::TursoDatabase>,
323}
324
325impl Debug for Database {
326    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
327        f.debug_struct("Database").finish()
328    }
329}
330
331impl Database {
332    /// Connect to the database.
333    pub fn connect(&self) -> Result<Connection> {
334        let conn = self.inner.connect()?;
335        Ok(Connection::create(conn, None))
336    }
337}
338
339/// A prepared statement.
340#[derive(Clone)]
341pub struct Statement {
342    conn: Connection,
343    inner: Arc<Mutex<Box<crate::turso_sdk_kit::rsapi::TursoStatement>>>,
344}
345
346struct Execute {
347    stmt: Statement,
348}
349
350assert_send_sync!(Execute);
351
352impl Future for Execute {
353    type Output = Result<u64>;
354
355    fn poll(
356        self: std::pin::Pin<&mut Self>,
357        cx: &mut std::task::Context<'_>,
358    ) -> std::task::Poll<Self::Output> {
359        match self.stmt.step(None, cx)? {
360            Poll::Ready(_) => {
361                let n_change = self.stmt.inner.lock().unwrap().n_change();
362                Poll::Ready(Ok(n_change as u64))
363            }
364            Poll::Pending => Poll::Pending,
365        }
366    }
367}
368
369impl Statement {
370    fn step(
371        &self,
372        columns: Option<usize>,
373        cx: &mut std::task::Context<'_>,
374    ) -> Poll<Result<Option<Row>>> {
375        let mut stmt = self.inner.lock().unwrap();
376        match stmt.step(Some(cx.waker()))? {
377            crate::turso_sdk_kit::rsapi::TursoStatusCode::Row => {
378                if let Some(columns) = columns {
379                    let mut values = Vec::with_capacity(columns);
380                    for i in 0..columns {
381                        let value = stmt.row_value(i)?;
382                        values.push(value);
383                    }
384                    Poll::Ready(Ok(Some(Row { values })))
385                } else {
386                    Poll::Ready(Err(Error::Misuse(
387                        "unexpected row during execution".to_string(),
388                    )))
389                }
390            }
391            crate::turso_sdk_kit::rsapi::TursoStatusCode::Done => Poll::Ready(Ok(None)),
392            crate::turso_sdk_kit::rsapi::TursoStatusCode::Io => {
393                stmt.run_io()?;
394                if let Some(extra_io) = &self.conn.extra_io {
395                    extra_io(cx.waker().clone())?;
396                }
397                Poll::Pending
398            }
399        }
400    }
401    /// Query the database with this prepared statement.
402    pub async fn query(&mut self, params: impl IntoParams) -> Result<Rows> {
403        self.reset()?;
404
405        let mut stmt = self.inner.lock().unwrap();
406        let params = params.into_params()?;
407        match params {
408            params::Params::None => (),
409            params::Params::Positional(values) => {
410                for (i, value) in values.into_iter().enumerate() {
411                    stmt.bind_positional(i + 1, value.into())?;
412                }
413            }
414            params::Params::Named(values) => {
415                for (name, value) in values.into_iter() {
416                    let position = stmt.named_position(name)?;
417                    stmt.bind_positional(position, value.into())?;
418                }
419            }
420        }
421        let rows = Rows::new(self.clone());
422        Ok(rows)
423    }
424
425    /// Execute this prepared statement.
426    pub async fn execute(&mut self, params: impl IntoParams) -> Result<u64> {
427        {
428            // Reset the statement before executing
429            self.inner.lock().unwrap().reset()?;
430        }
431        let params = params.into_params()?;
432        match params {
433            params::Params::None => (),
434            params::Params::Positional(values) => {
435                for (i, value) in values.into_iter().enumerate() {
436                    let mut stmt = self.inner.lock().unwrap();
437                    stmt.bind_positional(i + 1, value.into())?;
438                }
439            }
440            params::Params::Named(values) => {
441                for (name, value) in values.into_iter() {
442                    let mut stmt = self.inner.lock().unwrap();
443                    let position = stmt.named_position(name)?;
444                    stmt.bind_positional(position, value.into())?;
445                }
446            }
447        }
448
449        let execute = Execute { stmt: self.clone() };
450        execute.await
451    }
452
453    /// Returns the number of columns in the result set.
454    pub fn column_count(&self) -> usize {
455        self.inner.lock().unwrap().column_count()
456    }
457
458    /// Returns the name of the column at the given index.
459    pub fn column_name(&self, idx: usize) -> Result<String> {
460        let stmt = self.inner.lock().unwrap();
461        if idx >= stmt.column_count() {
462            return Err(Error::Misuse(format!(
463                "column index {idx} out of bounds (statement has {} columns)",
464                stmt.column_count()
465            )));
466        }
467        Ok(stmt
468            .column_name(idx)
469            .expect("column index must be within valid range"))
470    }
471
472    /// Returns the names of all columns in the result set.
473    pub fn column_names(&self) -> Vec<String> {
474        let stmt = self.inner.lock().unwrap();
475        let n = stmt.column_count();
476        (0..n)
477            .map(|i| {
478                stmt.column_name(i)
479                    .expect("column index must be within valid range")
480            })
481            .collect()
482    }
483
484    /// Returns the index of the column with the given name.
485    pub fn column_index(&self, name: &str) -> Result<usize> {
486        let stmt = self.inner.lock().unwrap();
487        let n = stmt.column_count();
488        for i in 0..n {
489            let col_name = stmt
490                .column_name(i)
491                .expect("column index must be within valid range");
492            if col_name.as_str().eq_ignore_ascii_case(name) {
493                return Ok(i);
494            }
495        }
496        Err(Error::Misuse(format!(
497            "column '{name}' not found in result set"
498        )))
499    }
500
501    /// Returns columns of the result of this prepared statement.
502    pub fn columns(&self) -> Vec<Column> {
503        let stmt = self.inner.lock().unwrap();
504
505        let n = stmt.column_count();
506
507        let mut cols = Vec::with_capacity(n);
508
509        for i in 0..n {
510            let name = stmt
511                .column_name(i)
512                .expect("column index must be within valid range");
513            let decl_type = stmt.column_decltype(i);
514            cols.push(Column { name, decl_type });
515        }
516
517        cols
518    }
519
520    /// Reset internal statement state after previous execution so it can be reused again
521    pub fn reset(&self) -> Result<()> {
522        let mut stmt = self.inner.lock().unwrap();
523        stmt.reset()?;
524        Ok(())
525    }
526
527    /// Returns the number of rows modified (insert/delete operations) by the most recent executed statement.
528    pub fn n_change(&self) -> u64 {
529        self.inner.lock().unwrap().n_change() as u64
530    }
531
532    /// Execute a query that returns the first [`Row`].
533    ///
534    /// # Errors
535    ///
536    /// - Returns `QueryReturnedNoRows` if no rows were returned.
537    pub async fn query_row(&mut self, params: impl IntoParams) -> Result<Row> {
538        let mut rows = self.query(params).await?;
539
540        let first_row = rows.next().await?.ok_or(Error::QueryReturnedNoRows)?;
541        // Discard remaining rows so that the statement is executed to completion
542        // Otherwise Drop of the statement will cause transaction rollback
543        while rows.next().await?.is_some() {}
544        Ok(first_row)
545    }
546}
547
548/// Column information.
549pub struct Column {
550    name: String,
551    decl_type: Option<String>,
552}
553
554impl Column {
555    /// Return the name of the column.
556    pub fn name(&self) -> &str {
557        &self.name
558    }
559
560    /// Returns the type of the column.
561    pub fn decl_type(&self) -> Option<&str> {
562        self.decl_type.as_deref()
563    }
564}
565
566pub trait IntoValue {
567    fn into_value(self) -> Result<Value>;
568}
569
570#[derive(Debug, Clone)]
571pub enum Params {
572    None,
573    Positional(Vec<Value>),
574    Named(Vec<(String, Value)>),
575}
576
577pub struct Transaction {}
578
579#[cfg(clt_turso_tests)]
580mod tests {
581    use super::*;
582    use tempfile::NamedTempFile;
583
584    #[tokio::test]
585    async fn test_database_persistence() -> Result<()> {
586        let temp_file = NamedTempFile::new().unwrap();
587        let db_path = temp_file.path().to_str().unwrap();
588
589        // First, create the database, a table, and insert some data
590        {
591            let db = Builder::new_local(db_path).build().await?;
592            let conn = db.connect()?;
593            conn.execute(
594                "CREATE TABLE test_persistence (id INTEGER PRIMARY KEY, name TEXT NOT NULL);",
595                (),
596            )
597            .await?;
598            conn.execute("INSERT INTO test_persistence (name) VALUES ('Alice');", ())
599                .await?;
600            conn.execute("INSERT INTO test_persistence (name) VALUES ('Bob');", ())
601                .await?;
602        } // db and conn are dropped here, simulating closing
603
604        // Now, re-open the database and check if the data is still there
605        let db = Builder::new_local(db_path).build().await?;
606        let conn = db.connect()?;
607
608        let mut rows = conn
609            .query("SELECT name FROM test_persistence ORDER BY id;", ())
610            .await?;
611
612        let row1 = rows.next().await?.expect("Expected first row");
613        assert_eq!(row1.get_value(0)?, Value::Text("Alice".to_string()));
614
615        let row2 = rows.next().await?.expect("Expected second row");
616        assert_eq!(row2.get_value(0)?, Value::Text("Bob".to_string()));
617
618        assert!(rows.next().await?.is_none(), "Expected no more rows");
619
620        Ok(())
621    }
622
623    #[tokio::test]
624    async fn test_database_persistence_many_frames() -> Result<()> {
625        let temp_file = NamedTempFile::new().unwrap();
626        let db_path = temp_file.path().to_str().unwrap();
627
628        const NUM_INSERTS: usize = 100;
629        const TARGET_STRING_LEN: usize = 1024; // 1KB
630
631        let mut original_data = Vec::with_capacity(NUM_INSERTS);
632        for i in 0..NUM_INSERTS {
633            let prefix = format!("test_string_{i:04}_");
634            let padding_len = TARGET_STRING_LEN.saturating_sub(prefix.len());
635            let padding: String = "A".repeat(padding_len);
636            original_data.push(format!("{prefix}{padding}"));
637        }
638
639        // First, create the database, a table, and insert many large strings
640        {
641            let db = Builder::new_local(db_path).build().await?;
642            let conn = db.connect()?;
643            conn.execute(
644                "CREATE TABLE test_large_persistence (id INTEGER PRIMARY KEY AUTOINCREMENT, data TEXT NOT NULL);",
645                (),
646            )
647            .await?;
648
649            for data_val in &original_data {
650                conn.execute(
651                    "INSERT INTO test_large_persistence (data) VALUES (?);",
652                    params::Params::Positional(vec![Value::Text(data_val.clone())]),
653                )
654                .await?;
655            }
656        } // db and conn are dropped here, simulating closing
657
658        {
659            // Now, re-open the database and check if the data is still there
660            let db = Builder::new_local(db_path).build().await?;
661            let conn = db.connect()?;
662
663            let mut rows = conn
664                .query("SELECT data FROM test_large_persistence ORDER BY id;", ())
665                .await?;
666
667            for (i, value) in original_data.iter().enumerate().take(NUM_INSERTS) {
668                let row = rows
669                    .next()
670                    .await?
671                    .unwrap_or_else(|| panic!("Expected row {i} but found None"));
672                assert_eq!(
673                    row.get_value(0)?,
674                    Value::Text(value.clone()),
675                    "Mismatch in retrieved data for row {i}"
676                );
677            }
678
679            assert!(
680                rows.next().await?.is_none(),
681                "Expected no more rows after retrieving all inserted data"
682            );
683
684            // Delete the WAL file only and try to re-open and query
685            let wal_path = format!("{db_path}-wal");
686            std::fs::remove_file(&wal_path)
687                .map_err(|e| eprintln!("Warning: Failed to delete WAL file for test: {e}"))
688                .unwrap();
689        }
690
691        // Attempt to re-open the database after deleting WAL and assert that table is missing.
692        let db_after_wal_delete = Builder::new_local(db_path).build().await?;
693        let conn_after_wal_delete = db_after_wal_delete.connect()?;
694
695        let query_result_after_wal_delete = conn_after_wal_delete
696            .query("SELECT data FROM test_large_persistence ORDER BY id;", ())
697            .await;
698
699        match query_result_after_wal_delete {
700            Ok(_) => panic!("Query succeeded after WAL deletion and DB reopen, but was expected to fail because the table definition should have been in the WAL."),
701            Err(Error::Error(msg)) => {
702                assert!(
703                    msg.contains("no such table: test_large_persistence"),
704                    "Expected 'test_large_persistence not found' error, but got: {msg}"
705                );
706            }
707            Err(e) => panic!(
708                "Expected SqlExecutionFailure for 'no such table', but got a different error: {e:?}"
709            ),
710        }
711
712        Ok(())
713    }
714
715    #[tokio::test]
716    async fn test_rows_column_names() -> Result<()> {
717        let db = Builder::new_local(":memory:").build().await?;
718        let conn = db.connect()?;
719        conn.execute(
720            "CREATE TABLE users (id INTEGER PRIMARY KEY, name TEXT, email TEXT);",
721            (),
722        )
723        .await?;
724        conn.execute(
725            "INSERT INTO users (name, email) VALUES ('Alice', 'alice@example.org');",
726            (),
727        )
728        .await?;
729
730        let rows = conn.query("SELECT id, name, email FROM users;", ()).await?;
731
732        // columns()
733        let columns = rows.columns();
734        let names: Vec<&str> = columns.iter().map(|c| c.name()).collect();
735        assert_eq!(names, vec!["id", "name", "email"]);
736
737        // column_count()
738        assert_eq!(rows.column_count(), 3);
739
740        // column_name()
741        assert_eq!(rows.column_name(0)?, "id");
742        assert_eq!(rows.column_name(1)?, "name");
743        assert_eq!(rows.column_name(2)?, "email");
744        assert!(rows.column_name(3).is_err());
745
746        // column_names()
747        assert_eq!(rows.column_names(), vec!["id", "name", "email"]);
748
749        // column_index()
750        assert_eq!(rows.column_index("id")?, 0);
751        assert_eq!(rows.column_index("name")?, 1);
752        assert_eq!(rows.column_index("email")?, 2);
753        assert_eq!(rows.column_index("EMAIL")?, 2); // case-insensitive
754        assert!(rows.column_index("nonexistent").is_err());
755
756        Ok(())
757    }
758
759    #[tokio::test]
760    async fn test_database_persistence_write_one_frame_many_times() -> Result<()> {
761        let temp_file = NamedTempFile::new().unwrap();
762        let db_path = temp_file.path().to_str().unwrap();
763
764        for i in 0..100 {
765            {
766                let db = Builder::new_local(db_path).build().await?;
767                let conn = db.connect()?;
768
769                conn.execute("CREATE TABLE IF NOT EXISTS test_persistence (id INTEGER PRIMARY KEY, name TEXT NOT NULL);", ()).await?;
770                conn.execute("INSERT INTO test_persistence (name) VALUES ('Alice');", ())
771                    .await?;
772            }
773            {
774                let db = Builder::new_local(db_path).build().await?;
775                let conn = db.connect()?;
776
777                let mut rows_iter = conn
778                    .query("SELECT count(*) FROM test_persistence;", ())
779                    .await?;
780                let rows = rows_iter.next().await?.unwrap();
781                assert_eq!(rows.get_value(0)?, Value::Integer(i as i64 + 1));
782                assert!(rows_iter.next().await?.is_none());
783            }
784        }
785
786        Ok(())
787    }
788
789    #[tokio::test]
790    async fn test_parallel_writes_and_wal_size() -> Result<()> {
791        let temp_dir = tempfile::tempdir().unwrap();
792        let db_path = temp_dir.path().join("test.db");
793        let db_path_str = db_path.to_str().unwrap();
794
795        let db = Builder::new_local(db_path_str).build().await?;
796        let conn = db.connect()?;
797        conn.execute(
798            "CREATE TABLE test_data (id INTEGER PRIMARY KEY AUTOINCREMENT, payload TEXT NOT NULL);",
799            (),
800        )
801        .await?;
802
803        // Generate a ~200KB payload
804        let payload = "X".repeat(200 * 1024);
805
806        // Parallel writes: spawn 8 connections, each inserting 5 rows
807        let mut handles = Vec::new();
808        for conn_id in 0..8u32 {
809            let db = db.clone();
810            let payload = payload.clone();
811            handles.push(tokio::spawn(async move {
812                let conn = db.connect().unwrap();
813                for row_id in 0..5u32 {
814                    let tag = format!("conn{conn_id}_row{row_id}");
815                    let data = format!("{tag}_{payload}");
816                    loop {
817                        match conn
818                            .execute(
819                                "INSERT INTO test_data (payload) VALUES (?);",
820                                params::Params::Positional(vec![Value::Text(data.clone())]),
821                            )
822                            .await
823                        {
824                            Ok(_) => break,
825                            Err(Error::Busy(_)) => {
826                                tokio::time::sleep(std::time::Duration::from_millis(10)).await;
827                                continue;
828                            }
829                            Err(e) => panic!("Insert failed: {e:?}"),
830                        }
831                    }
832                }
833            }));
834        }
835        for h in handles {
836            h.await.unwrap();
837        }
838
839        // Sequential writes: 3 more large inserts
840        for i in 0..3 {
841            let data = format!("sequential_{i}_{payload}");
842            conn.execute(
843                "INSERT INTO test_data (payload) VALUES (?);",
844                params::Params::Positional(vec![Value::Text(data)]),
845            )
846            .await?;
847        }
848
849        // Verify row count: 8*5 + 3 = 43
850        let mut rows = conn.query("SELECT count(*) FROM test_data;", ()).await?;
851        let row = rows.next().await?.unwrap();
852        assert_eq!(row.get_value(0)?, Value::Integer(43));
853
854        // Report WAL size
855        let wal_path = format!("{db_path_str}-wal");
856        let wal_size = std::fs::metadata(&wal_path).map(|m| m.len()).unwrap_or(0);
857        eprintln!(
858            "WAL size after all writes: {} bytes ({:.2} KB)",
859            wal_size,
860            wal_size as f64 / 1024.0
861        );
862        assert!(wal_size > 0, "WAL file should exist and be non-empty");
863
864        Ok(())
865    }
866}
867
868pub use crate::{named_params, params};