Skip to main content

clt_database/turso_sdk_kit/src/
rsapi.rs

1use std::{
2    collections::HashMap,
3    ffi::CString,
4    fmt::Display,
5    ops::Deref,
6    sync::{
7        atomic::{AtomicBool, AtomicUsize, Ordering},
8        Arc, Mutex, Once, RwLock, Weak,
9    },
10    task::Waker,
11    time::Duration,
12};
13
14use tracing::level_filters::LevelFilter;
15use tracing_subscriber::{
16    fmt::{self, format::Writer},
17    layer::{Context, SubscriberExt},
18    util::SubscriberInitExt,
19    EnvFilter, Layer,
20};
21use turso_core::{
22    storage::database::DatabaseFile, types::AsValueRef, Connection, Database, DatabaseOpts,
23    DatabaseStorage, EncryptionKey, IOResult, LimboError, OpenDbAsyncState, OpenFlags, QueryMode,
24    Statement, StepResult, IO,
25};
26
27use crate::turso_sdk_kit::{
28    assert_send, assert_sync,
29    capi::{self, c},
30    ConcurrentGuard,
31};
32
33assert_send!(TursoDatabase, TursoConnection, TursoStatement);
34assert_sync!(TursoDatabase);
35
36#[derive(Default)]
37pub struct SyncBusyGate {
38    active: AtomicBool,
39}
40
41impl SyncBusyGate {
42    pub fn set_active(&self, active: bool) {
43        self.active.store(active, Ordering::Release);
44    }
45
46    pub fn is_active(&self) -> bool {
47        self.active.load(Ordering::Acquire)
48    }
49}
50
51pub struct SyncBusyGuard {
52    gate: Arc<SyncBusyGate>,
53}
54
55impl SyncBusyGuard {
56    pub fn new(gate: Arc<SyncBusyGate>) -> Self {
57        gate.set_active(true);
58        Self { gate }
59    }
60}
61
62impl Drop for SyncBusyGuard {
63    fn drop(&mut self) {
64        self.gate.set_active(false);
65    }
66}
67
68pub use turso_core::types::FromValue;
69pub use turso_ext::{
70    AggCtx, ContextDestructor, FinalizeFunction, InitAggFunction, ResultCode, ScalarFunction,
71    StepFunction, Value as ExtensionValue, ValueDestructor,
72};
73pub type EncryptionOpts = turso_core::EncryptionOpts;
74pub type Value = turso_core::Value;
75pub type ValueRef<'a> = turso_core::types::ValueRef<'a>;
76pub type Text = turso_core::types::Text;
77pub type TextRef<'a> = turso_core::types::TextRef<'a>;
78pub type Numeric = turso_core::Numeric;
79pub type NonNan = turso_core::NonNan;
80
81pub struct TursoLog<'a> {
82    pub message: &'a str,
83    pub target: &'a str,
84    pub file: &'a str,
85    pub timestamp: u64,
86    pub line: usize,
87    pub level: &'a str,
88}
89
90type Logger = dyn Fn(TursoLog) + Send + Sync + 'static;
91pub struct TursoSetupConfig {
92    pub logger: Option<Box<Logger>>,
93    pub log_level: Option<String>,
94}
95
96fn logger_wrap(log: TursoLog<'_>, logger: unsafe extern "C" fn(*const c::turso_log_t)) {
97    let Ok(message_cstr) = std::ffi::CString::new(log.message) else {
98        return;
99    };
100    let Ok(target_cstr) = std::ffi::CString::new(log.target) else {
101        return;
102    };
103    let Ok(file_cstr) = std::ffi::CString::new(log.file) else {
104        return;
105    };
106    unsafe {
107        logger(&c::turso_log_t {
108            message: message_cstr.as_ptr(),
109            target: target_cstr.as_ptr(),
110            file: file_cstr.as_ptr(),
111            timestamp: log.timestamp,
112            line: log.line,
113            level: match log.level {
114                "TRACE" => capi::c::turso_tracing_level_t::TURSO_TRACING_LEVEL_TRACE,
115                "DEBUG" => capi::c::turso_tracing_level_t::TURSO_TRACING_LEVEL_DEBUG,
116                "INFO" => capi::c::turso_tracing_level_t::TURSO_TRACING_LEVEL_INFO,
117                "WARN" => capi::c::turso_tracing_level_t::TURSO_TRACING_LEVEL_WARN,
118                _ => capi::c::turso_tracing_level_t::TURSO_TRACING_LEVEL_ERROR,
119            },
120        })
121    };
122}
123
124impl TursoSetupConfig {
125    /// helper method to restore [TursoSetupConfig] instance from C representation
126    /// this method is used in the capi wrappers
127    ///
128    /// # Safety
129    /// [c::turso_config_t::log_level] field must be valid C-string pointer or null
130    pub unsafe fn from_capi(config: *const c::turso_config_t) -> Result<Self, TursoError> {
131        if config.is_null() {
132            return Err(TursoError::Misuse(
133                "config pointer must be not null".to_string(),
134            ));
135        }
136        let config = *config;
137        Ok(Self {
138            log_level: if !config.log_level.is_null() {
139                Some(str_from_c_str(config.log_level)?.to_string())
140            } else {
141                None
142            },
143            logger: if let Some(logger) = config.logger {
144                Some(Box::new(move |log| logger_wrap(log, logger)))
145            } else {
146                None
147            },
148        })
149    }
150}
151
152#[derive(Clone)]
153pub struct TursoDatabaseConfig {
154    /// path to the database file or ":memory:" for in-memory connection
155    pub path: String,
156
157    /// comma-separated list of experimental features to enable
158    /// this field is intentionally just a string in order to make enablement of experimental features as flexible as possible
159    pub experimental_features: Option<String>,
160
161    /// if true, library methods will return Io status code and delegate Io loop to the caller
162    /// if false, library will spin IO itself in case of Io status code and never return it to the caller
163    pub async_io: bool,
164
165    /// optional encryption parameters for local data encryption
166    /// as encryption is experimental - [Self::experimental_features] must have "encryption" in the list
167    pub encryption: Option<EncryptionOpts>,
168
169    /// optional VFS parameter explicitly specifying FS backend for the database.
170    /// Available options are:
171    /// - "memory": in-memory backend
172    /// - "syscall": generic syscall backend
173    /// - "io_uring": IO uring (supported only on Linux)
174    pub vfs: Option<String>,
175
176    /// optional custom IO provided by the caller
177    pub io: Option<Arc<dyn IO>>,
178
179    /// optional custom DatabaseStorage provided by the caller
180    /// if provided, caller must guarantee that IO used by the TursoDatabase will be consistent with underlying DatabaseStorage IO
181    pub db_file: Option<Arc<dyn DatabaseStorage>>,
182}
183
184impl TursoDatabaseConfig {
185    /// Build the typed [`turso_core::DatabaseOpts`] from the comma-separated
186    /// [`Self::experimental_features`] string. The feature-name tokens are the
187    /// SDK/CLI-facing names; unknown names are ignored and `"strict"` is a
188    /// no-op (strict tables are always enabled). This keeps the
189    /// string -> typed-options translation in the SDK layer where the string
190    /// representation lives, instead of in `turso_core`.
191    pub fn database_opts(&self) -> DatabaseOpts {
192        let mut opts = DatabaseOpts::new();
193        let Some(experimental_features) = &self.experimental_features else {
194            return opts;
195        };
196        for feature in experimental_features.split(',').map(|s| s.trim()) {
197            opts = match feature {
198                "views" => opts.with_views(true),
199                "index_method" => opts.with_index_method(true),
200                "custom_types" => opts.with_custom_types(true),
201                "autovacuum" => opts.with_autovacuum(true),
202                "vacuum" => opts.with_vacuum(true),
203                "encryption" => opts.with_encryption(true),
204                "attach" => opts.with_attach(true),
205                "generated_columns" => opts.with_generated_columns(true),
206                "multiprocess_wal" => opts.with_multiprocess_wal(true),
207                "without_rowid" => opts.with_without_rowid(true),
208                "mvcc_passive_checkpoint" => opts.with_experimental_mvcc_passive_checkpoint(true),
209                // "strict" is always enabled, kept for backwards compatibility
210                _ => opts,
211            };
212        }
213        opts
214    }
215}
216
217pub fn turso_slice_from_bytes(bytes: &[u8]) -> capi::c::turso_slice_ref_t {
218    capi::c::turso_slice_ref_t {
219        ptr: bytes.as_ptr() as *const std::ffi::c_void,
220        len: bytes.len(),
221    }
222}
223
224pub fn turso_slice_null() -> capi::c::turso_slice_ref_t {
225    capi::c::turso_slice_ref_t {
226        ptr: std::ptr::null(),
227        len: 0,
228    }
229}
230
231/// # Safety
232/// ptr must be valid C-string pointer or null
233pub unsafe fn str_from_c_str<'a>(ptr: *const std::ffi::c_char) -> Result<&'a str, TursoError> {
234    if ptr.is_null() {
235        return Err(TursoError::Misuse(
236            "expected zero terminated c string, got null pointer".to_string(),
237        ));
238    }
239    let c_str = std::ffi::CStr::from_ptr(ptr);
240    match c_str.to_str() {
241        Ok(s) => Ok(s),
242        Err(err) => Err(TursoError::Misuse(format!(
243            "expected zero terminated c-string representing utf-8 value: {err}"
244        ))),
245    }
246}
247
248/// # Safety
249/// memory range [ptr..ptr + len) must be valid
250pub unsafe fn str_from_slice<'a>(
251    ptr: *const std::ffi::c_char,
252    len: usize,
253) -> Result<&'a str, TursoError> {
254    let slice = bytes_from_slice(ptr, len)?;
255    match std::str::from_utf8(slice) {
256        Ok(s) => Ok(s),
257        Err(err) => Err(TursoError::Misuse(format!(
258            "expected string slice representing utf-8 value: {err}"
259        ))),
260    }
261}
262
263/// # Safety
264/// memory range [ptr..ptr + len) must be valid
265pub unsafe fn bytes_from_slice<'a>(
266    ptr: *const std::ffi::c_char,
267    len: usize,
268) -> Result<&'a [u8], TursoError> {
269    if len == 0 {
270        return Ok(&[]);
271    }
272    if ptr.is_null() {
273        return Err(TursoError::Misuse(
274            "expected slice, got null pointer".to_string(),
275        ));
276    }
277    Ok(std::slice::from_raw_parts(ptr as *const u8, len))
278}
279
280/// SAFETY: slice must points to the valid memory
281pub fn bytes_from_turso_slice<'a>(
282    slice: capi::c::turso_slice_ref_t,
283) -> Result<&'a [u8], TursoError> {
284    if slice.ptr.is_null() {
285        return Err(TursoError::Misuse(
286            "expected slice representing utf-8 value, got null".to_string(),
287        ));
288    }
289    Ok(unsafe { std::slice::from_raw_parts(slice.ptr as *const u8, slice.len) })
290}
291
292/// SAFETY: slice must points to the valid memory
293pub fn str_from_turso_slice<'a>(slice: capi::c::turso_slice_ref_t) -> Result<&'a str, TursoError> {
294    if slice.ptr.is_null() {
295        return Err(TursoError::Misuse(
296            "expected slice representing utf-8 value, got null".to_string(),
297        ));
298    }
299    let s = unsafe { std::slice::from_raw_parts(slice.ptr as *const u8, slice.len) };
300    match std::str::from_utf8(s) {
301        Ok(s) => Ok(s),
302        Err(err) => Err(TursoError::Misuse(format!(
303            "expected slice representing utf-8 value: {err}"
304        ))),
305    }
306}
307
308impl TursoDatabaseConfig {
309    /// helper method to restore [TursoSetupConfig] instance from C representation
310    /// this method is used in the capi wrappers
311    ///
312    /// # Safety
313    /// [c::turso_database_config_t::path] field must be valid C-string pointer
314    /// [c::turso_database_config_t::experimental_features] field must be valid C-string pointer or null
315    pub unsafe fn from_capi(config: *const c::turso_database_config_t) -> Result<Self, TursoError> {
316        if config.is_null() {
317            return Err(TursoError::Misuse(
318                "config pointer must be not null".to_string(),
319            ));
320        }
321        let config = *config;
322        let encryption_cipher = if !config.encryption_cipher.is_null() {
323            Some(str_from_c_str(config.encryption_cipher)?.to_string())
324        } else {
325            None
326        };
327        let encryption_hexkey = if !config.encryption_hexkey.is_null() {
328            Some(str_from_c_str(config.encryption_hexkey)?.to_string())
329        } else {
330            None
331        };
332        if encryption_cipher.is_some() != encryption_hexkey.is_some() {
333            return Err(TursoError::Misuse(
334                "either both encryption cipher and key must be set or no".to_string(),
335            ));
336        }
337        Ok(Self {
338            path: str_from_c_str(config.path)?.to_string(),
339            experimental_features: if !config.experimental_features.is_null() {
340                Some(str_from_c_str(config.experimental_features)?.to_string())
341            } else {
342                None
343            },
344            async_io: config.async_io != 0,
345            encryption: encryption_cipher.map(|encryption_cipher| EncryptionOpts {
346                cipher: encryption_cipher,
347                hexkey: encryption_hexkey.unwrap(),
348            }),
349            vfs: if !config.vfs.is_null() {
350                Some(str_from_c_str(config.vfs)?.to_string())
351            } else {
352                None
353            },
354            io: None,
355            db_file: None,
356        })
357    }
358}
359
360pub struct TursoDatabase {
361    config: TursoDatabaseConfig,
362    open_state: Mutex<TursoDatabaseOpenState>,
363    db: Arc<Mutex<Option<Arc<Database>>>>,
364    io: Mutex<Option<Arc<dyn turso_core::IO>>>,
365}
366
367/// Phase tracking for async TursoDatabase opening
368#[derive(Default, Clone, Copy)]
369pub enum TursoDatabaseOpenPhase {
370    #[default]
371    Init,
372    Opening,
373    Done,
374}
375
376/// State machine for async TursoDatabase opening
377pub struct TursoDatabaseOpenState {
378    phase: TursoDatabaseOpenPhase,
379    io: Option<Arc<dyn IO>>,
380    db_file: Option<Arc<dyn DatabaseStorage>>,
381    opts: Option<DatabaseOpts>,
382    open_flags: OpenFlags,
383    open_db_state: OpenDbAsyncState,
384}
385
386impl Default for TursoDatabaseOpenState {
387    fn default() -> Self {
388        Self::new()
389    }
390}
391
392impl TursoDatabaseOpenState {
393    pub fn new() -> Self {
394        Self {
395            phase: TursoDatabaseOpenPhase::Init,
396            io: None,
397            db_file: None,
398            opts: None,
399            open_flags: OpenFlags::default(),
400            open_db_state: OpenDbAsyncState::new(),
401        }
402    }
403}
404
405#[derive(Debug, Clone, Copy, PartialEq, Eq)]
406#[repr(u32)]
407pub enum TursoStatusCode {
408    Done,
409    Row,
410    Io,
411}
412
413#[derive(Debug, Clone)]
414pub enum TursoError {
415    Busy(String),
416    BusySnapshot(String),
417    Interrupt(String),
418    Error(String),
419    Misuse(String),
420    Constraint(String),
421    Readonly(String),
422    DatabaseFull(String),
423    NotAdb(String),
424    Corrupt(String),
425    IoError(std::io::ErrorKind, &'static str),
426}
427
428impl TursoStatusCode {
429    pub fn to_capi(self) -> capi::c::turso_status_code_t {
430        match self {
431            TursoStatusCode::Done => capi::c::turso_status_code_t::TURSO_DONE,
432            TursoStatusCode::Row => capi::c::turso_status_code_t::TURSO_ROW,
433            TursoStatusCode::Io => capi::c::turso_status_code_t::TURSO_IO,
434        }
435    }
436}
437
438fn result_code_to_result(result: ResultCode, operation: &str) -> Result<(), TursoError> {
439    if result.is_ok() {
440        Ok(())
441    } else {
442        Err(TursoError::Error(format!("{operation} failed: {result}")))
443    }
444}
445
446impl TursoError {
447    /// # Safety
448    /// error_opt_out must be a valid pointer or null
449    pub unsafe fn to_capi(
450        &self,
451        error_opt_out: *mut *const std::ffi::c_char,
452    ) -> capi::c::turso_status_code_t {
453        if !error_opt_out.is_null() {
454            let message = str_to_c_string(&self.to_string());
455            unsafe { *error_opt_out = message };
456        }
457        self.to_capi_code()
458    }
459    pub fn to_capi_code(&self) -> capi::c::turso_status_code_t {
460        match self {
461            TursoError::Busy(_) => capi::c::turso_status_code_t::TURSO_BUSY,
462            TursoError::BusySnapshot(_) => capi::c::turso_status_code_t::TURSO_BUSY_SNAPSHOT,
463            TursoError::Interrupt(_) => capi::c::turso_status_code_t::TURSO_INTERRUPT,
464            TursoError::Error(_) => capi::c::turso_status_code_t::TURSO_ERROR,
465            TursoError::Misuse(_) => capi::c::turso_status_code_t::TURSO_MISUSE,
466            TursoError::Constraint(_) => capi::c::turso_status_code_t::TURSO_CONSTRAINT,
467            TursoError::Readonly(_) => capi::c::turso_status_code_t::TURSO_READONLY,
468            TursoError::DatabaseFull(_) => capi::c::turso_status_code_t::TURSO_DATABASE_FULL,
469            TursoError::NotAdb(_) => capi::c::turso_status_code_t::TURSO_NOTADB,
470            TursoError::Corrupt(_) => capi::c::turso_status_code_t::TURSO_CORRUPT,
471            TursoError::IoError(..) => capi::c::turso_status_code_t::TURSO_IOERR,
472        }
473    }
474}
475
476impl Display for TursoError {
477    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
478        match self {
479            TursoError::Busy(s)
480            | TursoError::BusySnapshot(s)
481            | TursoError::Interrupt(s)
482            | TursoError::Error(s)
483            | TursoError::Misuse(s)
484            | TursoError::Constraint(s)
485            | TursoError::Readonly(s)
486            | TursoError::DatabaseFull(s)
487            | TursoError::NotAdb(s)
488            | TursoError::Corrupt(s) => f.write_str(s),
489            TursoError::IoError(kind, op) => write!(f, "I/O error ({op}): {kind}"),
490        }
491    }
492}
493
494pub fn str_to_c_string(message: &str) -> *const std::ffi::c_char {
495    let Ok(message) = std::ffi::CString::new(message) else {
496        return std::ptr::null();
497    };
498    message.into_raw()
499}
500
501pub fn c_string_to_str(ptr: *const std::ffi::c_char) -> std::ffi::CString {
502    unsafe { std::ffi::CString::from_raw(ptr as *mut std::ffi::c_char) }
503}
504
505impl From<LimboError> for TursoError {
506    fn from(value: LimboError) -> Self {
507        match value {
508            LimboError::ForeignKeyConstraint(e) | LimboError::Constraint(e) => {
509                TursoError::Constraint(e)
510            }
511            LimboError::Corrupt(e) => TursoError::Corrupt(e),
512            LimboError::NotADB => TursoError::NotAdb("file is not a database".to_string()),
513            LimboError::DatabaseFull(e) => TursoError::DatabaseFull(e),
514            LimboError::ReadOnly => TursoError::Readonly("database is readonly".to_string()),
515            LimboError::Busy => TursoError::Busy("database is locked".to_string()),
516            // Same-connection rejections carry SQLITE_BUSY semantics, but the
517            // caller must finish/reset its own statement rather than wait.
518            err @ LimboError::StatementsInProgress(_) => TursoError::Busy(err.to_string()),
519            LimboError::BusySnapshot => TursoError::BusySnapshot(
520                "database snapshot is stale, rollback and retry the transaction".to_string(),
521            ),
522            LimboError::CompletionError(turso_core::CompletionError::IOError(kind, op)) => {
523                TursoError::IoError(kind, op)
524            }
525            _ => TursoError::Error(value.to_string()),
526        }
527    }
528}
529
530fn sync_busy_error() -> TursoError {
531    TursoError::Busy("database is locked".to_string())
532}
533
534fn sync_operation_active(sync_busy: Option<&Arc<SyncBusyGate>>) -> bool {
535    sync_busy.is_some_and(|gate| gate.is_active())
536}
537
538fn map_sync_transient_error(
539    sync_busy: Option<&Arc<SyncBusyGate>>,
540    error: TursoError,
541) -> TursoError {
542    if sync_busy.is_none() {
543        return error;
544    }
545    match error {
546        TursoError::Error(message)
547            if message == "Database schema changed"
548                || message.starts_with("I/O error: short read on page") =>
549        {
550            sync_busy_error()
551        }
552        other => other,
553    }
554}
555
556static LOGGER: RwLock<Option<Box<Logger>>> = RwLock::new(None);
557static SETUP: Once = Once::new();
558
559struct CallbackLayer<F>
560where
561    F: Fn(TursoLog) + Send + Sync + 'static,
562{
563    callback: F,
564}
565
566impl<S, F> tracing_subscriber::Layer<S> for CallbackLayer<F>
567where
568    S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>,
569    F: Fn(TursoLog) + Send + Sync + 'static,
570{
571    fn on_event(&self, event: &tracing::Event<'_>, _ctx: Context<'_, S>) {
572        let mut buffer = String::new();
573        let mut visitor = fmt::format::DefaultVisitor::new(Writer::new(&mut buffer), true);
574
575        event.record(&mut visitor);
576
577        let log = TursoLog {
578            level: event.metadata().level().as_str(),
579            target: event.metadata().target(),
580            message: &buffer,
581            timestamp: std::time::SystemTime::now()
582                .duration_since(std::time::UNIX_EPOCH)
583                .map(|t| t.as_secs())
584                .unwrap_or(0),
585            file: event.metadata().file().unwrap_or(""),
586            line: event.metadata().line().unwrap_or(0) as usize,
587        };
588
589        (self.callback)(log);
590    }
591}
592
593pub fn turso_setup(config: TursoSetupConfig) -> Result<(), TursoError> {
594    fn callback(log: TursoLog<'_>) {
595        let Ok(logger) = LOGGER.try_read() else {
596            return;
597        };
598
599        if let Some(logger) = logger.as_ref() {
600            logger(log)
601        }
602    }
603
604    if let Some(logger) = config.logger {
605        let mut guard = LOGGER.write().unwrap();
606        *guard = Some(logger);
607    }
608
609    let level_filter = if let Some(log_level) = &config.log_level {
610        match log_level.as_ref() {
611            "error" => Some(LevelFilter::ERROR),
612            "warn" => Some(LevelFilter::WARN),
613            "info" => Some(LevelFilter::INFO),
614            "debug" => Some(LevelFilter::DEBUG),
615            "trace" => Some(LevelFilter::TRACE),
616            _ => return Err(TursoError::Error("unknown log level".to_string())),
617        }
618    } else {
619        None
620    };
621
622    SETUP.call_once(|| {
623        if let Some(level_filter) = level_filter {
624            tracing_subscriber::registry()
625                .with(CallbackLayer { callback }.with_filter(level_filter))
626                .init();
627        } else {
628            tracing_subscriber::registry()
629                .with(CallbackLayer { callback }.with_filter(EnvFilter::from_default_env()))
630                .init();
631        }
632    });
633
634    Ok(())
635}
636
637impl TursoDatabase {
638    /// return turso version
639    pub const fn version() -> &'static str {
640        "0.7.2-clt.1"
641    }
642    /// method to get [turso_core::Database] instance which can be useful for code which integrates with sdk-kit
643    pub fn db_core(&self) -> Result<Arc<turso_core::Database>, TursoError> {
644        let db = self.db.lock().unwrap();
645        match &*db {
646            Some(db) => Ok(db.clone()),
647            None => Err(TursoError::Misuse("database must be opened".to_string())),
648        }
649    }
650
651    /// method to get [turso_core::IO] instance which can be useful for code which integrates with sdk-kit
652    pub fn io(&self) -> Result<Arc<dyn turso_core::IO>, TursoError> {
653        let io = self.io.lock().unwrap();
654        match &*io {
655            Some(io) => Ok(io.clone()),
656            None => Err(TursoError::Misuse("io must be opened".to_string())),
657        }
658    }
659
660    /// create database holder struct but do not initialize it yet
661    /// this can be useful for some environments, where IO operations must be executed in certain fashion (and open do IO under the hood)
662    pub fn new(config: TursoDatabaseConfig) -> Arc<Self> {
663        Arc::new(Self {
664            config,
665            db: Arc::new(Mutex::new(None)),
666            open_state: Mutex::new(TursoDatabaseOpenState::new()),
667            io: Mutex::new(None),
668        })
669    }
670
671    /// Get the config IO or open a new vfs IO from the config
672    fn open_vfs_io(&self) -> Result<Arc<dyn turso_core::IO>, TursoError> {
673        let io: Arc<dyn turso_core::IO + 'static> = if let Some(io) = &self.config.io {
674            io.clone()
675        } else {
676            match self.config.vfs.as_deref() {
677                Some("memory") => Arc::new(turso_core::MemoryIO::new()),
678                Some("syscall") => {
679                    #[cfg(all(target_family = "unix", not(miri)))]
680                    {
681                        Arc::new(turso_core::UnixIO::new().map_err(|e| {
682                            TursoError::Error(format!(
683                                "unable to create generic syscall backend: {e}"
684                            ))
685                        })?)
686                    }
687                    #[cfg(any(not(target_family = "unix"), miri))]
688                    {
689                        Arc::new(turso_core::PlatformIO::new().map_err(|e| {
690                            TursoError::Error(format!(
691                                "unable to create generic syscall backend: {e}"
692                            ))
693                        })?)
694                    }
695                }
696                #[cfg(all(target_os = "linux", not(miri)))]
697                Some("io_uring") => Arc::new(turso_core::UringIO::new().map_err(|e| {
698                    TursoError::Error(format!("unable to create io_uring backend: {e}"))
699                })?),
700                #[cfg(all(target_os = "windows", not(miri)))]
701                Some("experimental_win_iocp") => {
702                    Arc::new(turso_core::WindowsIOCP::new().map_err(|e| {
703                        TursoError::Error(format!("unable to create win_iocp backend: {e}"))
704                    })?)
705                }
706                #[cfg(any(not(target_os = "linux"), miri))]
707                Some("io_uring") => {
708                    return Err(TursoError::Error(
709                        "io_uring is only available on Linux targets".to_string(),
710                    ));
711                }
712                #[cfg(any(not(target_os = "windows"), miri))]
713                Some("experimental_win_iocp") => {
714                    return Err(TursoError::Error(
715                        "win_iocp is only available on Windows targets".to_string(),
716                    ));
717                }
718                Some(vfs) => {
719                    Database::io_for_vfs(vfs).map_err(|e| TursoError::Error(format!("{e}")))?
720                }
721                None => match self.config.path.as_str() {
722                    ":memory:" => Arc::new(turso_core::MemoryIO::new()),
723                    _ => Arc::new(turso_core::PlatformIO::new()?),
724                },
725            }
726        };
727        Ok(io)
728    }
729
730    /// Async version of database opening that returns IOResult.
731    /// Caller must drive the IO loop and pass state between calls.
732    /// This is useful for environments where IO operations must be executed in a specific fashion.
733    pub fn open(&self) -> Result<IOResult<()>, TursoError> {
734        loop {
735            let mut state = self.open_state.lock().unwrap();
736            match state.phase {
737                TursoDatabaseOpenPhase::Init => {
738                    let inner_db = self.db.lock().unwrap();
739                    if inner_db.is_some() {
740                        return Err(TursoError::Misuse(
741                            "database must be opened only once".to_string(),
742                        ));
743                    }
744                    // keep lock for the whole method since open_async must be called only once and never will be called concurrently
745
746                    let io: Arc<dyn turso_core::IO> = self.open_vfs_io()?;
747
748                    // Store the IO so that it can be retrieved with `io()` call even if the database is still opening
749                    *self.io.lock().unwrap() = Some(io.clone());
750
751                    // Opts must be computed BEFORE the file open so we can apply
752                    // OpenFlags::NoLock when multiprocess WAL is enabled — taking
753                    // the OS-level fcntl lock here would block every other
754                    // multiprocess process from opening the same file.
755                    let opts = self.config.database_opts();
756
757                    if self.config.encryption.is_some() && !opts.enable_encryption {
758                        return Err(TursoError::Error(
759                            "encryption is experimental and must be explicitly enabled through experimental features list".to_string(),
760                        ));
761                    }
762
763                    let mut open_flags = OpenFlags::default();
764                    if opts.enable_multiprocess_wal {
765                        open_flags |= OpenFlags::NoLock;
766                    }
767                    let db_file = if let Some(db_file) = &self.config.db_file {
768                        db_file.clone()
769                    } else {
770                        let file = io.open_file(&self.config.path, open_flags, true)?;
771                        Arc::new(DatabaseFile::new(file))
772                    };
773
774                    state.io = Some(io);
775                    state.db_file = Some(db_file);
776                    state.opts = Some(opts);
777                    state.open_flags = open_flags;
778                    state.phase = TursoDatabaseOpenPhase::Opening;
779                }
780
781                TursoDatabaseOpenPhase::Opening => {
782                    let io = state
783                        .io
784                        .as_ref()
785                        .expect("io must be initialized in Init phase")
786                        .clone();
787                    let db_file = state
788                        .db_file
789                        .as_ref()
790                        .expect("db_file must be initialized in Init phase")
791                        .clone();
792                    let opts = state.opts.expect("opts must be initialized in Init phase");
793                    let open_flags = state.open_flags;
794
795                    match Database::open_with_flags_async(
796                        &mut state.open_db_state,
797                        io.clone(),
798                        &self.config.path,
799                        db_file,
800                        open_flags,
801                        opts,
802                        self.config.encryption.clone(),
803                        None,
804                    )? {
805                        IOResult::Done(db) => {
806                            let mut inner_db = self.db.lock().unwrap();
807                            *inner_db = Some(db);
808                            state.phase = TursoDatabaseOpenPhase::Done;
809                            return Ok(IOResult::Done(()));
810                        }
811                        IOResult::IO(io_completion) => {
812                            if self.config.async_io {
813                                return Ok(IOResult::IO(io_completion));
814                            } else {
815                                io_completion.wait(io.deref())?;
816                            }
817                        }
818                    }
819                }
820
821                TursoDatabaseOpenPhase::Done => {
822                    return Ok(IOResult::Done(()));
823                }
824            }
825        }
826    }
827
828    /// creates database connection
829    /// database must be already opened with [Self::open] method
830    pub fn connect(&self) -> Result<Arc<TursoConnection>, TursoError> {
831        let inner_db = self.db.lock().unwrap();
832        let Some(db) = inner_db.as_ref() else {
833            return Err(TursoError::Misuse(
834                "database must be opened first".to_string(),
835            ));
836        };
837
838        // Parse encryption key if configured - needed for connect_with_encryption
839        // which sets up encryption context before reading pages
840        let encryption_key = if let Some(ref encryption_opts) = self.config.encryption {
841            Some(EncryptionKey::from_hex_string(&encryption_opts.hexkey)?)
842        } else {
843            None
844        };
845
846        // Use connect_with_encryption to properly set up encryption context
847        // before the pager reads page 1. This is required for encrypted databases.
848        let connection = db.connect_with_encryption(encryption_key)?;
849
850        Ok(TursoConnection::new(&self.config, connection))
851    }
852
853    /// helper method to get C raw container with TursoDatabase instance
854    /// this method is used in the capi wrappers
855    pub fn to_capi(self: Arc<Self>) -> *mut capi::c::turso_database_t {
856        Arc::into_raw(self) as *mut capi::c::turso_database_t
857    }
858
859    /// helper method to restore TursoDatabase ref from C raw container
860    /// this method is used in the capi wrappers
861    ///
862    /// # Safety
863    /// value must be a pointer returned from [Self::to_capi] method
864    pub unsafe fn ref_from_capi<'a>(
865        value: *const capi::c::turso_database_t,
866    ) -> Result<&'a Self, TursoError> {
867        if value.is_null() {
868            Err(TursoError::Misuse("got null pointer".to_string()))
869        } else {
870            Ok(&*(value as *const Self))
871        }
872    }
873
874    /// helper method to restore TursoDatabase instance from C raw container
875    /// this method is used in the capi wrappers
876    ///
877    /// # Safety
878    /// value must be a pointer returned from [Self::to_capi] method
879    pub unsafe fn arc_from_capi(value: *const capi::c::turso_database_t) -> Arc<Self> {
880        Arc::from_raw(value as *const Self)
881    }
882}
883
884struct CachedStatement {
885    program: Arc<turso_core::PreparedProgram>,
886    query_mode: QueryMode,
887}
888
889#[derive(Clone)]
890pub struct TursoConnection {
891    async_io: bool,
892    concurrent_guard: Arc<ConcurrentGuard>,
893    connection: Arc<Connection>,
894    sync_busy: Option<Arc<SyncBusyGate>>,
895    cached_statements: Arc<Mutex<HashMap<String, Arc<CachedStatement>>>>,
896    /// Weak refs to every statement handle created by this connection, keyed
897    /// by a monotonic ID. Statements remove themselves on drop, so this map
898    /// only ever contains live entries. `close()` upgrades each remaining
899    /// handle and sets it to `None` to release `Arc<Connection>` → `Arc<Database>`.
900    stmts: StmtRegistry,
901    next_stmt_id: Arc<AtomicUsize>,
902}
903
904impl TursoConnection {
905    pub fn new(config: &TursoDatabaseConfig, connection: Arc<Connection>) -> Arc<Self> {
906        Self::new_with_sync_busy(config, connection, None)
907    }
908
909    pub fn new_with_sync_busy(
910        config: &TursoDatabaseConfig,
911        connection: Arc<Connection>,
912        sync_busy: Option<Arc<SyncBusyGate>>,
913    ) -> Arc<Self> {
914        Arc::new(Self {
915            async_io: config.async_io,
916            connection,
917            sync_busy,
918            concurrent_guard: Arc::new(ConcurrentGuard::new()),
919            cached_statements: Arc::new(Mutex::new(HashMap::new())),
920            stmts: Arc::new(Mutex::new(HashMap::new())),
921            next_stmt_id: Arc::new(AtomicUsize::new(0)),
922        })
923    }
924
925    fn sync_operation_active(&self) -> bool {
926        sync_operation_active(self.sync_busy.as_ref())
927    }
928
929    fn map_sync_transient_error(&self, error: TursoError) -> TursoError {
930        map_sync_transient_error(self.sync_busy.as_ref(), error)
931    }
932    /// Set busy timeout for the connection
933    pub fn set_busy_timeout(&self, duration: Duration) {
934        self.connection.set_busy_timeout(duration);
935    }
936    /// Request interruption of the statement currently running on this connection.
937    /// Mirrors `sqlite3_interrupt`: the in-flight `step`/`execute` aborts with an
938    /// `Interrupt` error. Safe to call from another thread. If no statement is
939    /// active the request is ignored.
940    pub fn interrupt(&self) {
941        self.connection.interrupt();
942    }
943    /// Set the maximum wall-clock duration a single statement is allowed to run
944    /// before it is interrupted. `Duration::ZERO` disables the timeout.
945    pub fn set_query_timeout(&self, duration: Duration) {
946        self.connection.set_query_timeout(duration);
947    }
948    /// Get the current per-statement query timeout (`Duration::ZERO` when disabled).
949    pub fn get_query_timeout(&self) -> Duration {
950        self.connection.get_query_timeout()
951    }
952    pub fn get_auto_commit(&self) -> bool {
953        self.connection.get_auto_commit()
954    }
955    pub fn last_insert_rowid(&self) -> i64 {
956        self.connection.last_insert_rowid()
957    }
958
959    #[allow(clippy::too_many_arguments)]
960    pub fn register_external_scalar_function(
961        &self,
962        name: String,
963        argc: i32,
964        deterministic: bool,
965        context: usize,
966        callback: ScalarFunction,
967        context_destructor: Option<ContextDestructor>,
968        value_destructor: Option<ValueDestructor>,
969    ) -> Result<(), TursoError> {
970        let name = CString::new(name).map_err(|err| {
971            TursoError::Misuse(format!(
972                "external scalar function name contains interior NUL: {err}"
973            ))
974        })?;
975        let api = unsafe { self.connection._build_turso_ext() };
976        let result = unsafe {
977            (api.register_scalar_function)(
978                api.ctx,
979                name.as_ptr(),
980                argc,
981                deterministic,
982                context,
983                callback,
984                context_destructor,
985                value_destructor,
986            )
987        };
988        unsafe { self.connection._free_extension_ctx(api) };
989        result_code_to_result(result, "register external scalar function")
990    }
991
992    #[allow(clippy::too_many_arguments)]
993    pub fn register_external_aggregate_function(
994        &self,
995        name: String,
996        argc: i32,
997        context: usize,
998        init: InitAggFunction,
999        step: StepFunction,
1000        finalize: FinalizeFunction,
1001        context_destructor: Option<ContextDestructor>,
1002        aggregate_destructor: Option<ContextDestructor>,
1003        value_destructor: Option<ValueDestructor>,
1004    ) -> Result<(), TursoError> {
1005        let name = CString::new(name).map_err(|err| {
1006            TursoError::Misuse(format!(
1007                "external aggregate function name contains interior NUL: {err}"
1008            ))
1009        })?;
1010        let api = unsafe { self.connection._build_turso_ext() };
1011        let result = unsafe {
1012            (api.register_aggregate_function)(
1013                api.ctx,
1014                name.as_ptr(),
1015                argc,
1016                context,
1017                init,
1018                step,
1019                finalize,
1020                context_destructor,
1021                aggregate_destructor,
1022                value_destructor,
1023            )
1024        };
1025        unsafe { self.connection._free_extension_ctx(api) };
1026        result_code_to_result(result, "register external aggregate function")
1027    }
1028
1029    pub fn unregister_external_function(&self, name: &str) -> Result<(), TursoError> {
1030        let name = CString::new(name).map_err(|err| {
1031            TursoError::Misuse(format!(
1032                "external function name contains interior NUL: {err}"
1033            ))
1034        })?;
1035        let api = unsafe { self.connection._build_turso_ext() };
1036        let result = unsafe { (api.unregister_function)(api.ctx, name.as_ptr()) };
1037        unsafe { self.connection._free_extension_ctx(api) };
1038        result_code_to_result(result, "unregister external function")
1039    }
1040
1041    pub fn register_external_collation(
1042        &self,
1043        name: String,
1044        context: usize,
1045        callback: turso_core::ContextCollationFunction,
1046        context_destructor: Option<ContextDestructor>,
1047    ) {
1048        self.connection
1049            .register_external_collation(name, context, callback, context_destructor);
1050    }
1051
1052    pub fn unregister_external_collation(&self, name: &str) {
1053        self.connection.unregister_external_collation(name);
1054    }
1055
1056    pub fn set_load_extension_enabled(&self, enabled: bool) {
1057        self.connection.set_load_extension_enabled(enabled);
1058    }
1059
1060    pub fn load_extension(&self, path: &str) -> Result<(), TursoError> {
1061        turso_core::resolve_ext_path(path)
1062            .and_then(|path| self.connection.load_extension(path))
1063            .map_err(TursoError::from)
1064    }
1065
1066    /// prepares single SQL statement
1067    pub fn prepare_single(&self, sql: impl AsRef<str>) -> Result<Box<TursoStatement>, TursoError> {
1068        if self.sync_operation_active() {
1069            return Err(sync_busy_error());
1070        }
1071        let statement = self
1072            .connection
1073            .prepare(sql)
1074            .map_err(TursoError::from)
1075            .map_err(|error| self.map_sync_transient_error(error))?;
1076        let handle: StatementHandle = Arc::new(Mutex::new(Some(statement)));
1077        let stmt_id = self.track_stmt(&handle);
1078        Ok(Box::new(TursoStatement {
1079            concurrent_guard: self.concurrent_guard.clone(),
1080            async_io: self.async_io,
1081            sync_busy: self.sync_busy.clone(),
1082            handle,
1083            stmt_id,
1084            stmts: self.stmts.clone(),
1085        }))
1086    }
1087
1088    /// Prepare a statement from the provided SQL string and cache it for future use.
1089    pub fn prepare_cached(&self, sql: impl AsRef<str>) -> Result<Box<TursoStatement>, TursoError> {
1090        if self.sync_operation_active() {
1091            return Err(sync_busy_error());
1092        }
1093        let sql_str = sql.as_ref();
1094
1095        // Check if we have a cached version
1096        if let Some(cached) = self.cached_statements.lock().unwrap().get(sql_str) {
1097            if cached.program.is_compatible_with(&self.connection) {
1098                let program = turso_core::Program::from_prepared(
1099                    cached.program.clone(),
1100                    self.connection.clone(),
1101                );
1102                let statement =
1103                    Statement::new(program, self.connection.get_pager(), cached.query_mode, 0);
1104                let handle: StatementHandle = Arc::new(Mutex::new(Some(statement)));
1105                let stmt_id = self.track_stmt(&handle);
1106                return Ok(Box::new(TursoStatement {
1107                    concurrent_guard: self.concurrent_guard.clone(),
1108                    async_io: self.async_io,
1109                    sync_busy: self.sync_busy.clone(),
1110                    handle,
1111                    stmt_id,
1112                    stmts: self.stmts.clone(),
1113                }));
1114            }
1115        }
1116
1117        // Not cached, prepare it fresh
1118        let statement = self
1119            .connection
1120            .prepare(sql_str)
1121            .map_err(TursoError::from)
1122            .map_err(|error| self.map_sync_transient_error(error))?;
1123
1124        // Cache it for future use
1125        let cached = Arc::new(CachedStatement {
1126            program: statement.get_program().prepared().clone(),
1127            query_mode: statement.get_query_mode(),
1128        });
1129        self.cached_statements
1130            .lock()
1131            .unwrap()
1132            .insert(sql_str.to_string(), cached);
1133
1134        let handle: StatementHandle = Arc::new(Mutex::new(Some(statement)));
1135        let stmt_id = self.track_stmt(&handle);
1136        Ok(Box::new(TursoStatement {
1137            concurrent_guard: self.concurrent_guard.clone(),
1138            async_io: self.async_io,
1139            sync_busy: self.sync_busy.clone(),
1140            handle,
1141            stmt_id,
1142            stmts: self.stmts.clone(),
1143        }))
1144    }
1145
1146    /// prepares first SQL statement from the string and return prepared statement and position after the end of the parsed statement
1147    /// this method can be useful if SDK provides an execute(...) method which run all statements from the provided input in sequence
1148    pub fn prepare_first(
1149        &self,
1150        sql: impl AsRef<str>,
1151    ) -> Result<Option<(Box<TursoStatement>, usize)>, TursoError> {
1152        if self.sync_operation_active() {
1153            return Err(sync_busy_error());
1154        }
1155        match self
1156            .connection
1157            .consume_stmt(sql)
1158            .map_err(TursoError::from)
1159            .map_err(|error| self.map_sync_transient_error(error))?
1160        {
1161            Some((statement, position)) => {
1162                let handle: StatementHandle = Arc::new(Mutex::new(Some(statement)));
1163                let stmt_id = self.track_stmt(&handle);
1164                Ok(Some((
1165                    Box::new(TursoStatement {
1166                        async_io: self.async_io,
1167                        concurrent_guard: Arc::new(ConcurrentGuard::new()),
1168                        sync_busy: self.sync_busy.clone(),
1169                        handle,
1170                        stmt_id,
1171                        stmts: self.stmts.clone(),
1172                    }),
1173                    position,
1174                )))
1175            }
1176            None => Ok(None),
1177        }
1178    }
1179
1180    /// close the connection preventing any further operations executed over it
1181    /// SAFETY: caller must guarantee that no ongoing operations are running over connection before calling close(...) method
1182    pub fn close(&self) -> Result<(), TursoError> {
1183        // Finalize all outstanding statements to release their Arc chain:
1184        // Statement → Program → Arc<Connection> → Arc<Database>.
1185        // Without this, un-finalized statements keep the Database alive in
1186        // DATABASE_MANAGER, causing stale databases after file renames.
1187        let mut stmts = self.stmts.lock().unwrap();
1188        for (_id, weak) in stmts.drain() {
1189            if let Some(handle) = weak.upgrade() {
1190                // Setting to None drops the turso_core::Statement,
1191                // releasing Arc<Connection> → Arc<Database>.
1192                *handle.lock().unwrap() = None;
1193            }
1194        }
1195        self.connection.close()?;
1196        Ok(())
1197    }
1198
1199    /// low-level method used only by the Rust SDK
1200    pub fn cacheflush(&self) -> Result<(), TursoError> {
1201        let completions = self.connection.cacheflush()?;
1202        let pager = self.connection.get_pager();
1203        for c in completions {
1204            pager.io.wait_for_completion(c)?;
1205        }
1206        Ok(())
1207    }
1208
1209    /// helper method to get C raw container to the TursoConnection instance
1210    /// this method is used in the capi wrappers
1211    pub fn to_capi(self: Arc<Self>) -> *mut capi::c::turso_connection_t {
1212        Arc::into_raw(self) as *mut capi::c::turso_connection_t
1213    }
1214
1215    /// helper method to restore TursoConnection ref from C raw container
1216    /// this method is used in the capi wrappers
1217    ///
1218    /// # Safety
1219    /// value must be a pointer returned from [Self::to_capi] method
1220    pub unsafe fn ref_from_capi<'a>(
1221        value: *const capi::c::turso_connection_t,
1222    ) -> Result<&'a Self, TursoError> {
1223        if value.is_null() {
1224            Err(TursoError::Misuse("got null pointer".to_string()))
1225        } else {
1226            Ok(&*(value as *const Self))
1227        }
1228    }
1229
1230    /// helper method to restore TursoConnection instance from C raw container
1231    /// this method is used in the capi wrappers
1232    ///
1233    /// # Safety
1234    /// value must be a pointer returned from [Self::to_capi] method
1235    pub unsafe fn arc_from_capi(value: *const capi::c::turso_connection_t) -> Arc<Self> {
1236        Arc::from_raw(value as *const Self)
1237    }
1238
1239    /// Register a statement handle and return its ID. The statement removes
1240    /// itself from the registry on drop via its `stmt_id` + `stmts` ref.
1241    fn track_stmt(&self, handle: &StatementHandle) -> usize {
1242        let id = self.next_stmt_id.fetch_add(1, Ordering::Relaxed);
1243        self.stmts
1244            .lock()
1245            .unwrap()
1246            .insert(id, Arc::downgrade(handle));
1247        id
1248    }
1249}
1250
1251/// Shared ownership of a `turso_core::Statement` that can be explicitly finalized.
1252/// When the inner `Option` is set to `None`, the statement is considered finalized
1253/// and all operations on it will return errors / defaults.
1254pub(crate) type StatementHandle = Arc<Mutex<Option<Statement>>>;
1255type StmtRegistry = Arc<Mutex<HashMap<usize, Weak<Mutex<Option<Statement>>>>>>;
1256
1257const FINALIZED_ERR: &str = "statement has been finalized";
1258
1259/// Advance one step of a statement's execution.
1260/// Factored out of `TursoStatement` so it can be called while holding
1261/// the `StatementHandle` lock without re-entrancy issues.
1262fn step_inner(
1263    stmt: &mut Statement,
1264    async_io: bool,
1265    waker: Option<&Waker>,
1266) -> Result<TursoStatusCode, TursoError> {
1267    loop {
1268        let result = if let Some(waker) = waker {
1269            stmt.step_with_waker(waker)
1270        } else {
1271            stmt.step()
1272        };
1273        return match result? {
1274            StepResult::Done => Ok(TursoStatusCode::Done),
1275            StepResult::Row => Ok(TursoStatusCode::Row),
1276            StepResult::Busy => Err(TursoError::Busy("database is locked".to_string())),
1277            StepResult::Interrupt => Err(TursoError::Interrupt("interrupted".to_string())),
1278            StepResult::IO | StepResult::Yield => {
1279                if async_io {
1280                    Ok(TursoStatusCode::Io)
1281                } else {
1282                    stmt._io().step()?;
1283                    continue;
1284                }
1285            }
1286        };
1287    }
1288}
1289
1290pub struct TursoStatement {
1291    async_io: bool,
1292    concurrent_guard: Arc<ConcurrentGuard>,
1293    sync_busy: Option<Arc<SyncBusyGate>>,
1294    pub(crate) handle: StatementHandle,
1295    stmt_id: usize,
1296    stmts: StmtRegistry,
1297}
1298
1299impl Drop for TursoStatement {
1300    fn drop(&mut self) {
1301        self.stmts.lock().unwrap().remove(&self.stmt_id);
1302    }
1303}
1304
1305#[derive(Debug, Clone)]
1306pub struct TursoExecutionResult {
1307    pub status: TursoStatusCode,
1308    pub rows_changed: u64,
1309}
1310
1311impl TursoStatement {
1312    /// return amount of row modifications (insert/delete operations) made by the most recent executed statement
1313    pub fn n_change(&self) -> i64 {
1314        let handle = self.handle.lock().unwrap();
1315        match handle.as_ref() {
1316            Some(stmt) => stmt.n_change(),
1317            None => 0,
1318        }
1319    }
1320    /// returns parameters count for the statement
1321    pub fn parameters_count(&self) -> usize {
1322        let handle = self.handle.lock().unwrap();
1323        match handle.as_ref() {
1324            Some(stmt) => stmt.parameters_count(),
1325            None => 0,
1326        }
1327    }
1328    /// Returns the name of the parameter at the given 1-based index,
1329    /// including its SQL prefix (e.g. `:name`, `@name`, `$name`).
1330    /// Returns None for positional-only (`?`) parameters or out-of-range indices.
1331    pub fn parameter_name(&self, index: usize) -> Option<String> {
1332        let handle = self.handle.lock().unwrap();
1333        let stmt = handle.as_ref()?;
1334        let index = index.try_into().ok()?;
1335        stmt.parameters().name(index)
1336    }
1337    /// binds positional parameter at the corresponding index (1-based)
1338    pub fn bind_positional(
1339        &mut self,
1340        index: usize,
1341        value: turso_core::Value,
1342    ) -> Result<(), TursoError> {
1343        let mut handle = self.handle.lock().unwrap();
1344        let stmt = handle
1345            .as_mut()
1346            .ok_or_else(|| TursoError::Misuse(FINALIZED_ERR.to_string()))?;
1347        let Ok(index) = index.try_into() else {
1348            return Err(TursoError::Misuse(
1349                "bind index must be non-zero".to_string(),
1350            ));
1351        };
1352        if !stmt.parameters().has_slot(index) {
1353            return Err(TursoError::Misuse(format!(
1354                "bind index {index} is out of bounds"
1355            )));
1356        }
1357        stmt.bind_at(index, value)?;
1358        Ok(())
1359    }
1360    /// named parameter position.
1361    ///
1362    /// The name must include the SQL placeholder prefix, e.g. `:name`, `@name`, `$name`, or `?1`.
1363    pub fn named_position(&mut self, name: impl AsRef<str>) -> Result<usize, TursoError> {
1364        let handle = self.handle.lock().unwrap();
1365        let stmt = handle
1366            .as_ref()
1367            .ok_or_else(|| TursoError::Misuse(FINALIZED_ERR.to_string()))?;
1368        let name = name.as_ref();
1369        if let Some(index) = stmt.parameter_index(name) {
1370            return Ok(index.into());
1371        }
1372
1373        if name.starts_with('?') {
1374            let maybe_index = name
1375                .strip_prefix('?')
1376                .and_then(|value| value.parse::<usize>().ok())
1377                .and_then(|value| value.try_into().ok());
1378            if let Some(index) = maybe_index {
1379                if stmt.parameters().is_indexed(index) {
1380                    return Ok(index.into());
1381                }
1382            }
1383        }
1384
1385        Err(TursoError::Error(format!(
1386            "named parameter {name} not found"
1387        )))
1388    }
1389    /// make one execution step of the statement
1390    /// method returns [TursoStatusCode::Done] if execution is finished
1391    /// method returns [TursoStatusCode::Row] if execution generated a row
1392    /// method returns [TursoStatusCode::Io] if async_io was set and execution needs IO in order to make progress
1393    #[inline]
1394    pub fn step(&mut self, waker: Option<&Waker>) -> Result<TursoStatusCode, TursoError> {
1395        if sync_operation_active(self.sync_busy.as_ref()) {
1396            return Err(sync_busy_error());
1397        }
1398        let guard = self.concurrent_guard.clone();
1399        let _guard = guard.try_use()?;
1400        let mut handle = self.handle.lock().unwrap();
1401        let stmt = handle
1402            .as_mut()
1403            .ok_or_else(|| TursoError::Misuse(FINALIZED_ERR.to_string()))?;
1404        step_inner(stmt, self.async_io, waker)
1405            .map_err(|error| map_sync_transient_error(self.sync_busy.as_ref(), error))
1406    }
1407
1408    /// execute statement to completion
1409    /// method returns [TursoStatusCode::Done] if execution completed
1410    /// method returns [TursoStatusCode::Io] if async_io was set and execution needs IO in order to make progress
1411    pub fn execute(&mut self, waker: Option<&Waker>) -> Result<TursoExecutionResult, TursoError> {
1412        if sync_operation_active(self.sync_busy.as_ref()) {
1413            return Err(sync_busy_error());
1414        }
1415        let guard = self.concurrent_guard.clone();
1416        let _guard = guard.try_use()?;
1417        let mut handle = self.handle.lock().unwrap();
1418        let stmt = handle
1419            .as_mut()
1420            .ok_or_else(|| TursoError::Misuse(FINALIZED_ERR.to_string()))?;
1421
1422        loop {
1423            let status = step_inner(stmt, self.async_io, waker)
1424                .map_err(|error| map_sync_transient_error(self.sync_busy.as_ref(), error))?;
1425            if status == TursoStatusCode::Row {
1426                continue;
1427            } else if status == TursoStatusCode::Io {
1428                return Ok(TursoExecutionResult {
1429                    status,
1430                    rows_changed: 0,
1431                });
1432            } else if status == TursoStatusCode::Done {
1433                return Ok(TursoExecutionResult {
1434                    status: TursoStatusCode::Done,
1435                    rows_changed: stmt.n_change() as u64,
1436                });
1437            }
1438            return Err(TursoError::Error(format!(
1439                "internal error: unexpected status code: {status:?}",
1440            )));
1441        }
1442    }
1443    /// run iteration of the IO backend
1444    pub fn run_io(&self) -> Result<(), TursoError> {
1445        let handle = self.handle.lock().unwrap();
1446        let stmt = handle
1447            .as_ref()
1448            .ok_or_else(|| TursoError::Misuse(FINALIZED_ERR.to_string()))?;
1449        stmt._io().step()?;
1450        Ok(())
1451    }
1452    /// get row value as an owned Value
1453    #[inline]
1454    pub fn row_value(&self, index: usize) -> Result<turso_core::Value, TursoError> {
1455        if sync_operation_active(self.sync_busy.as_ref()) {
1456            return Err(sync_busy_error());
1457        }
1458        let handle = self.handle.lock().unwrap();
1459        let stmt = handle
1460            .as_ref()
1461            .ok_or_else(|| TursoError::Misuse(FINALIZED_ERR.to_string()))?;
1462        let Some(row) = stmt.row() else {
1463            return Err(TursoError::Misuse("statement holds no row".to_string()));
1464        };
1465        if index >= row.len() {
1466            return Err(TursoError::Misuse(
1467                "attempt to access row value out of bounds".to_string(),
1468            ));
1469        }
1470        Ok(row.get_value(index).as_value_ref().to_owned())
1471    }
1472    /// returns column count
1473    pub fn column_count(&self) -> usize {
1474        let handle = self.handle.lock().unwrap();
1475        match handle.as_ref() {
1476            Some(stmt) => stmt.num_columns(),
1477            None => 0,
1478        }
1479    }
1480    /// returns column name
1481    pub fn column_name(&self, index: usize) -> Result<String, TursoError> {
1482        let handle = self.handle.lock().unwrap();
1483        let stmt = handle
1484            .as_ref()
1485            .ok_or_else(|| TursoError::Misuse(FINALIZED_ERR.to_string()))?;
1486        if index >= stmt.num_columns() {
1487            return Err(TursoError::Misuse("column index out of bounds".to_string()));
1488        }
1489        Ok(stmt.get_column_name(index).into_owned())
1490    }
1491    /// returns column declared type (e.g. "INTEGER", "TEXT", "DATETIME", etc.)
1492    pub fn column_decltype(&self, index: usize) -> Option<String> {
1493        let handle = self.handle.lock().unwrap();
1494        let stmt = handle.as_ref()?;
1495        if index >= stmt.num_columns() {
1496            return None;
1497        }
1498        stmt.get_column_decltype(index)
1499    }
1500
1501    /// Returns rich type information for the column at `index`.
1502    ///
1503    /// Wraps [`turso_core::Statement::get_column_type_info`]. Returns `None`
1504    /// when the statement has been finalized, when the index is out of
1505    /// bounds, when the connection does not have the experimental custom-
1506    /// types feature enabled (the underlying call errors and we surface that
1507    /// as "no info"; the C ABI has no error channel), when the statement is
1508    /// in EXPLAIN mode, or when the expression behind the column has no
1509    /// determined affinity.
1510    pub fn column_type_info(&self, index: usize) -> Option<turso_core::ColumnTypeInfo> {
1511        let handle = self.handle.lock().unwrap();
1512        let stmt = handle.as_ref()?;
1513        if index >= stmt.num_columns() {
1514            return None;
1515        }
1516        stmt.get_column_type_info(index).ok().flatten()
1517    }
1518    /// finalize statement execution
1519    /// this method must be called in the end of statement execution (either successfull or not)
1520    pub fn finalize(&mut self, waker: Option<&Waker>) -> Result<TursoStatusCode, TursoError> {
1521        let guard = self.concurrent_guard.clone();
1522        let _guard = guard.try_use()?;
1523        let mut handle = self.handle.lock().unwrap();
1524        if let Some(stmt) = handle.as_mut() {
1525            while stmt.execution_state().is_running() {
1526                let status = step_inner(stmt, self.async_io, waker)?;
1527                if status == TursoStatusCode::Io {
1528                    return Ok(status);
1529                }
1530            }
1531        }
1532        // Drop the inner statement to release the Arc chain
1533        *handle = None;
1534        Ok(TursoStatusCode::Done)
1535    }
1536    /// reset internal statement state and bindings
1537    pub fn reset(&mut self) -> Result<(), TursoError> {
1538        let mut handle = self.handle.lock().unwrap();
1539        let stmt = handle
1540            .as_mut()
1541            .ok_or_else(|| TursoError::Misuse(FINALIZED_ERR.to_string()))?;
1542        stmt.reset()?;
1543        stmt.clear_bindings();
1544        Ok(())
1545    }
1546
1547    /// helper method to get C raw container to the TursoStatement instance
1548    /// this method is used in the capi wrappers
1549    pub fn to_capi(self: Box<Self>) -> *mut capi::c::turso_statement_t {
1550        Box::into_raw(self) as *mut capi::c::turso_statement_t
1551    }
1552
1553    /// helper method to restore TursoStatement ref from C raw container
1554    /// this method is used in the capi wrappers
1555    ///
1556    /// # Safety
1557    /// value must be a pointer returned from [Self::to_capi] method
1558    pub unsafe fn ref_from_capi<'a>(
1559        value: *const capi::c::turso_statement_t,
1560    ) -> Result<&'a mut Self, TursoError> {
1561        if value.is_null() {
1562            Err(TursoError::Misuse("got null pointer".to_string()))
1563        } else {
1564            Ok(&mut *(value as *mut Self))
1565        }
1566    }
1567
1568    /// helper method to restore TursoStatement instance from C raw container
1569    /// this method is used in the capi wrappers
1570    ///
1571    /// # Safety
1572    /// value must be a pointer returned from [Self::to_capi] method
1573    pub unsafe fn box_from_capi(value: *const capi::c::turso_statement_t) -> Box<Self> {
1574        Box::from_raw(value as *mut Self)
1575    }
1576}
1577
1578#[cfg(clt_turso_tests)]
1579mod tests {
1580    use crate::turso_sdk_kit::rsapi::{
1581        TursoDatabase, TursoDatabaseConfig, TursoError, TursoStatusCode, FINALIZED_ERR,
1582    };
1583    use turso_core::Value;
1584
1585    fn config_with_features(features: Option<&str>) -> TursoDatabaseConfig {
1586        TursoDatabaseConfig {
1587            path: ":memory:".to_string(),
1588            experimental_features: features.map(str::to_string),
1589            async_io: false,
1590            encryption: None,
1591            vfs: None,
1592            io: None,
1593            db_file: None,
1594        }
1595    }
1596
1597    #[test]
1598    pub fn database_opts_maps_experimental_features() {
1599        // No features -> all defaults.
1600        assert_eq!(
1601            config_with_features(None).database_opts(),
1602            turso_core::DatabaseOpts::new()
1603        );
1604
1605        // Each token toggles its corresponding flag.
1606        let opts = config_with_features(Some(
1607            "views,index_method,custom_types,autovacuum,vacuum,encryption,attach,generated_columns,multiprocess_wal,without_rowid",
1608        ))
1609        .database_opts();
1610        assert!(opts.enable_views);
1611        assert!(opts.enable_index_method);
1612        assert!(opts.enable_custom_types);
1613        assert!(opts.enable_autovacuum);
1614        assert!(opts.enable_vacuum);
1615        assert!(opts.enable_encryption);
1616        assert!(opts.enable_attach);
1617        assert!(opts.enable_generated_columns);
1618        assert!(opts.enable_multiprocess_wal);
1619        assert!(opts.enable_without_rowid);
1620
1621        // Whitespace is trimmed; `strict` and unknown names are ignored.
1622        let opts = config_with_features(Some(" views , strict , unknown_one ")).database_opts();
1623        assert!(opts.enable_views);
1624        assert_eq!(
1625            config_with_features(Some("strict,unknown")).database_opts(),
1626            turso_core::DatabaseOpts::new()
1627        );
1628    }
1629
1630    #[test]
1631    pub fn test_db_concurrent_use() {
1632        use std::sync::{Arc, Barrier};
1633
1634        let mut errors = Vec::new();
1635        for _ in 0..16 {
1636            let db = TursoDatabase::new(TursoDatabaseConfig {
1637                path: ":memory:".to_string(),
1638                experimental_features: None,
1639                async_io: false,
1640                encryption: None,
1641                vfs: None,
1642                io: None,
1643                db_file: None,
1644            });
1645            let result = db.open().unwrap();
1646            assert!(!result.is_io());
1647            let conn = db.connect().unwrap();
1648            let stmt1 = conn
1649                .prepare_single("SELECT * FROM generate_series(1, 100000)")
1650                .unwrap();
1651            let stmt2 = conn
1652                .prepare_single("SELECT * FROM generate_series(1, 100000)")
1653                .unwrap();
1654
1655            // Use a barrier to ensure both threads start executing at the same time
1656            let barrier = Arc::new(Barrier::new(2));
1657            let mut threads = Vec::new();
1658            for mut stmt in [stmt1, stmt2] {
1659                let barrier_clone = Arc::clone(&barrier);
1660                let thread = std::thread::spawn(move || {
1661                    barrier_clone.wait();
1662                    stmt.execute(None)
1663                });
1664                threads.push(thread);
1665            }
1666            let mut results = Vec::new();
1667            for thread in threads {
1668                results.push(thread.join().unwrap());
1669            }
1670            assert!(
1671                !(results[0].is_err() && results[1].is_err()),
1672                "results: {results:?}",
1673            );
1674            if results[0].is_err() || results[1].is_err() {
1675                errors.push(
1676                    results[0]
1677                        .clone()
1678                        .err()
1679                        .or(results[1].clone().err())
1680                        .unwrap(),
1681                );
1682            }
1683        }
1684        println!("{errors:?}");
1685        assert!(
1686            !errors.is_empty(),
1687            "misuse errors should be very likely with the test setup: {errors:?}"
1688        );
1689        assert!(
1690            errors.iter().all(|e| matches!(e, TursoError::Misuse(_))),
1691            "all errors must have Misuse code: {errors:?}"
1692        );
1693    }
1694
1695    #[test]
1696    pub fn test_db_rsapi_use() {
1697        let db = TursoDatabase::new(TursoDatabaseConfig {
1698            path: ":memory:".to_string(),
1699            experimental_features: None,
1700            async_io: false,
1701            encryption: None,
1702            vfs: None,
1703            io: None,
1704            db_file: None,
1705        });
1706        let result = db.open().unwrap();
1707        assert!(!result.is_io());
1708        let conn = db.connect().unwrap();
1709        let mut stmt = conn
1710            .prepare_single("SELECT * FROM generate_series(1, 10000)")
1711            .unwrap();
1712        assert_eq!(stmt.execute(None).unwrap().status, TursoStatusCode::Done);
1713    }
1714
1715    #[test]
1716    pub fn test_named_position_requires_prefixed_name() {
1717        let db = TursoDatabase::new(TursoDatabaseConfig {
1718            path: ":memory:".to_string(),
1719            experimental_features: None,
1720            async_io: false,
1721            encryption: None,
1722            vfs: None,
1723            io: None,
1724            db_file: None,
1725        });
1726        let result = db.open().unwrap();
1727        assert!(!result.is_io());
1728
1729        let conn = db.connect().unwrap();
1730        let mut stmt = conn
1731            .prepare_single("SELECT :new_name, @other_name, $third_name")
1732            .unwrap();
1733
1734        assert_eq!(stmt.named_position(":new_name").unwrap(), 1);
1735        assert!(stmt.named_position("new_name").is_err());
1736        assert!(stmt.named_position("?1").is_err());
1737
1738        assert_eq!(stmt.named_position("@other_name").unwrap(), 2);
1739        assert!(stmt.named_position("other_name").is_err());
1740
1741        assert_eq!(stmt.named_position("$third_name").unwrap(), 3);
1742        assert!(stmt.named_position("third_name").is_err());
1743    }
1744
1745    #[test]
1746    pub fn test_bind_positional_rejects_out_of_bounds_index() {
1747        let db = TursoDatabase::new(TursoDatabaseConfig {
1748            path: ":memory:".to_string(),
1749            experimental_features: None,
1750            async_io: false,
1751            encryption: None,
1752            vfs: None,
1753            io: None,
1754            db_file: None,
1755        });
1756        let result = db.open().unwrap();
1757        assert!(!result.is_io());
1758
1759        let conn = db.connect().unwrap();
1760        let mut stmt = conn.prepare_single("SELECT ?1").unwrap();
1761
1762        stmt.bind_positional(1, Value::from_i64(42)).unwrap();
1763
1764        let err = stmt.bind_positional(2, Value::from_i64(7)).unwrap_err();
1765        assert!(matches!(err, TursoError::Misuse(_)));
1766    }
1767
1768    #[test]
1769    pub fn test_execute_update_with_prefixed_named_parameters() {
1770        let db = TursoDatabase::new(TursoDatabaseConfig {
1771            path: ":memory:".to_string(),
1772            experimental_features: None,
1773            async_io: false,
1774            encryption: None,
1775            vfs: None,
1776            io: None,
1777            db_file: None,
1778        });
1779        let result = db.open().unwrap();
1780        assert!(!result.is_io());
1781
1782        let conn = db.connect().unwrap();
1783
1784        let mut create_stmt = conn
1785            .prepare_single("CREATE TABLE simple (id INTEGER PRIMARY KEY, name TEXT NOT NULL)")
1786            .unwrap();
1787        assert_eq!(
1788            create_stmt.execute(None).unwrap().status,
1789            TursoStatusCode::Done
1790        );
1791
1792        let mut insert_stmt = conn
1793            .prepare_single("INSERT INTO simple (name) VALUES ('original_name')")
1794            .unwrap();
1795        assert_eq!(
1796            insert_stmt.execute(None).unwrap().status,
1797            TursoStatusCode::Done
1798        );
1799
1800        let mut update_stmt = conn
1801            .prepare_single("UPDATE simple SET name = :new_name WHERE name = :old_name")
1802            .unwrap();
1803
1804        let new_name_position = update_stmt.named_position(":new_name").unwrap();
1805        update_stmt
1806            .bind_positional(new_name_position, Value::build_text("updated_name"))
1807            .unwrap();
1808        let old_name_position = update_stmt.named_position(":old_name").unwrap();
1809        update_stmt
1810            .bind_positional(old_name_position, Value::build_text("original_name"))
1811            .unwrap();
1812
1813        let update_result = update_stmt.execute(None).unwrap();
1814        assert_eq!(update_result.status, TursoStatusCode::Done);
1815        assert_eq!(update_result.rows_changed, 1);
1816    }
1817
1818    #[test]
1819    pub fn test_execute_update_with_mixed_placeholders() {
1820        let db = TursoDatabase::new(TursoDatabaseConfig {
1821            path: ":memory:".to_string(),
1822            experimental_features: None,
1823            async_io: false,
1824            encryption: None,
1825            vfs: None,
1826            io: None,
1827            db_file: None,
1828        });
1829        let result = db.open().unwrap();
1830        assert!(!result.is_io());
1831
1832        let conn = db.connect().unwrap();
1833
1834        let mut create_stmt = conn
1835            .prepare_single("CREATE TABLE mixed (id INTEGER PRIMARY KEY, name TEXT NOT NULL, email TEXT, age INTEGER)")
1836            .unwrap();
1837        assert_eq!(
1838            create_stmt.execute(None).unwrap().status,
1839            TursoStatusCode::Done
1840        );
1841
1842        let mut insert_stmt = conn
1843            .prepare_single(
1844                "INSERT INTO mixed (name, email, age) VALUES ('alice', 'alice@old.com', 25)",
1845            )
1846            .unwrap();
1847        assert_eq!(
1848            insert_stmt.execute(None).unwrap().status,
1849            TursoStatusCode::Done
1850        );
1851
1852        let mut update_stmt = conn
1853            .prepare_single("UPDATE mixed SET email = ?, age = :new_age WHERE name = ?")
1854            .unwrap();
1855
1856        assert_eq!(update_stmt.named_position("?1").unwrap(), 1);
1857        assert_eq!(update_stmt.named_position(":new_age").unwrap(), 2);
1858        assert!(update_stmt.named_position("new_age").is_err());
1859        assert_eq!(update_stmt.named_position("?3").unwrap(), 3);
1860
1861        update_stmt
1862            .bind_positional(1, Value::build_text("alice@new.com"))
1863            .unwrap();
1864        let age_position = update_stmt.named_position(":new_age").unwrap();
1865        update_stmt
1866            .bind_positional(age_position, Value::from_i64(30))
1867            .unwrap();
1868        update_stmt
1869            .bind_positional(3, Value::build_text("alice"))
1870            .unwrap();
1871
1872        let update_result = update_stmt.execute(None).unwrap();
1873        assert_eq!(update_result.status, TursoStatusCode::Done);
1874        assert_eq!(update_result.rows_changed, 1);
1875    }
1876
1877    #[test]
1878    pub fn test_select_named_and_positional_mapping_stays_sql_order() {
1879        let db = TursoDatabase::new(TursoDatabaseConfig {
1880            path: ":memory:".to_string(),
1881            experimental_features: None,
1882            async_io: false,
1883            encryption: None,
1884            vfs: None,
1885            io: None,
1886            db_file: None,
1887        });
1888        let result = db.open().unwrap();
1889        assert!(!result.is_io());
1890
1891        let conn = db.connect().unwrap();
1892        let mut create_stmt = conn
1893            .prepare_single("CREATE TABLE simple (name TEXT NOT NULL)")
1894            .unwrap();
1895        assert_eq!(
1896            create_stmt.execute(None).unwrap().status,
1897            TursoStatusCode::Done
1898        );
1899
1900        let mut stmt = conn
1901            .prepare_single("SELECT :named FROM simple WHERE name = ?")
1902            .unwrap();
1903
1904        assert_eq!(stmt.named_position(":named").unwrap(), 1);
1905        assert!(stmt.named_position("named").is_err());
1906        assert_eq!(stmt.named_position("?2").unwrap(), 2);
1907    }
1908
1909    #[test]
1910    pub fn test_named_and_indexed_alias_share_slot() {
1911        let db = TursoDatabase::new(TursoDatabaseConfig {
1912            path: ":memory:".to_string(),
1913            experimental_features: None,
1914            async_io: false,
1915            encryption: None,
1916            vfs: None,
1917            io: None,
1918            db_file: None,
1919        });
1920        let result = db.open().unwrap();
1921        assert!(!result.is_io());
1922
1923        let conn = db.connect().unwrap();
1924        let mut stmt = conn
1925            .prepare_single("SELECT :v AS named_slot, ?1 AS pos_slot")
1926            .unwrap();
1927
1928        assert_eq!(stmt.named_position(":v").unwrap(), 1);
1929        assert!(stmt.named_position("v").is_err());
1930        assert!(stmt.named_position("?1").is_err());
1931
1932        stmt.bind_positional(1, Value::from_i64(7)).unwrap();
1933        assert_eq!(stmt.step(None).unwrap(), TursoStatusCode::Row);
1934        assert_eq!(stmt.row_value(0).unwrap().as_int(), Some(7));
1935        assert_eq!(stmt.row_value(1).unwrap().as_int(), Some(7));
1936    }
1937
1938    #[test]
1939    pub fn test_sparse_positional_index_uses_declared_slot() {
1940        let db = TursoDatabase::new(TursoDatabaseConfig {
1941            path: ":memory:".to_string(),
1942            experimental_features: None,
1943            async_io: false,
1944            encryption: None,
1945            vfs: None,
1946            io: None,
1947            db_file: None,
1948        });
1949        let result = db.open().unwrap();
1950        assert!(!result.is_io());
1951
1952        let conn = db.connect().unwrap();
1953        let mut stmt = conn.prepare_single("SELECT ?3").unwrap();
1954
1955        assert_eq!(stmt.parameters_count(), 3);
1956        stmt.bind_positional(1, Value::from_i64(1)).unwrap();
1957
1958        stmt.bind_positional(3, Value::from_i64(9)).unwrap();
1959        assert_eq!(stmt.step(None).unwrap(), TursoStatusCode::Row);
1960        assert_eq!(stmt.row_value(0).unwrap().as_int(), Some(9));
1961    }
1962
1963    #[test]
1964    pub fn test_sparse_positional_index_count_matches_sqlite() {
1965        let db = TursoDatabase::new(TursoDatabaseConfig {
1966            path: ":memory:".to_string(),
1967            experimental_features: None,
1968            async_io: false,
1969            encryption: None,
1970            vfs: None,
1971            io: None,
1972            db_file: None,
1973        });
1974        let result = db.open().unwrap();
1975        assert!(!result.is_io());
1976
1977        let conn = db.connect().unwrap();
1978        let mut stmt = conn.prepare_single("SELECT ?3").unwrap();
1979
1980        assert_eq!(stmt.parameters_count(), 3);
1981        assert!(stmt.named_position("?1").is_err());
1982        assert_eq!(stmt.named_position("?3").unwrap(), 3);
1983
1984        stmt.bind_positional(3, Value::from_i64(11)).unwrap();
1985        assert_eq!(stmt.step(None).unwrap(), TursoStatusCode::Row);
1986        assert_eq!(stmt.row_value(0).unwrap().as_int(), Some(11));
1987    }
1988
1989    #[test]
1990    pub fn test_insert_with_mixed_placeholders() {
1991        let db = TursoDatabase::new(TursoDatabaseConfig {
1992            path: ":memory:".to_string(),
1993            experimental_features: None,
1994            async_io: false,
1995            encryption: None,
1996            vfs: None,
1997            io: None,
1998            db_file: None,
1999        });
2000        let result = db.open().unwrap();
2001        assert!(!result.is_io());
2002
2003        let conn = db.connect().unwrap();
2004        let mut create_stmt = conn
2005            .prepare_single("CREATE TABLE users (name TEXT NOT NULL, age INTEGER NOT NULL)")
2006            .unwrap();
2007        assert_eq!(
2008            create_stmt.execute(None).unwrap().status,
2009            TursoStatusCode::Done
2010        );
2011
2012        let mut insert_stmt = conn
2013            .prepare_single("INSERT INTO users (name, age) VALUES (?, :age)")
2014            .unwrap();
2015
2016        assert_eq!(insert_stmt.named_position("?1").unwrap(), 1);
2017        assert_eq!(insert_stmt.named_position(":age").unwrap(), 2);
2018        assert!(insert_stmt.named_position("age").is_err());
2019
2020        insert_stmt
2021            .bind_positional(1, Value::build_text("alice"))
2022            .unwrap();
2023        let age_position = insert_stmt.named_position(":age").unwrap();
2024        insert_stmt
2025            .bind_positional(age_position, Value::from_i64(30))
2026            .unwrap();
2027
2028        let insert_result = insert_stmt.execute(None).unwrap();
2029        assert_eq!(insert_result.status, TursoStatusCode::Done);
2030        assert_eq!(insert_result.rows_changed, 1);
2031
2032        let mut verify_stmt = conn
2033            .prepare_single("SELECT age FROM users WHERE name = 'alice'")
2034            .unwrap();
2035        assert_eq!(verify_stmt.step(None).unwrap(), TursoStatusCode::Row);
2036        assert_eq!(verify_stmt.row_value(0).unwrap().as_int(), Some(30));
2037    }
2038
2039    #[test]
2040    pub fn test_delete_with_mixed_placeholders() {
2041        let db = TursoDatabase::new(TursoDatabaseConfig {
2042            path: ":memory:".to_string(),
2043            experimental_features: None,
2044            async_io: false,
2045            encryption: None,
2046            vfs: None,
2047            io: None,
2048            db_file: None,
2049        });
2050        let result = db.open().unwrap();
2051        assert!(!result.is_io());
2052
2053        let conn = db.connect().unwrap();
2054        let mut create_stmt = conn
2055            .prepare_single("CREATE TABLE users (name TEXT NOT NULL, age INTEGER NOT NULL)")
2056            .unwrap();
2057        assert_eq!(
2058            create_stmt.execute(None).unwrap().status,
2059            TursoStatusCode::Done
2060        );
2061
2062        let mut seed_stmt = conn
2063            .prepare_single("INSERT INTO users (name, age) VALUES ('alice', 30), ('bob', 40)")
2064            .unwrap();
2065        assert_eq!(
2066            seed_stmt.execute(None).unwrap().status,
2067            TursoStatusCode::Done
2068        );
2069
2070        let mut delete_stmt = conn
2071            .prepare_single("DELETE FROM users WHERE name = ? AND age = :age")
2072            .unwrap();
2073
2074        assert_eq!(delete_stmt.named_position("?1").unwrap(), 1);
2075        assert_eq!(delete_stmt.named_position(":age").unwrap(), 2);
2076        assert!(delete_stmt.named_position("age").is_err());
2077
2078        delete_stmt
2079            .bind_positional(1, Value::build_text("alice"))
2080            .unwrap();
2081        let age_position = delete_stmt.named_position(":age").unwrap();
2082        delete_stmt
2083            .bind_positional(age_position, Value::from_i64(30))
2084            .unwrap();
2085
2086        let delete_result = delete_stmt.execute(None).unwrap();
2087        assert_eq!(delete_result.status, TursoStatusCode::Done);
2088        assert_eq!(delete_result.rows_changed, 1);
2089
2090        let mut verify_stmt = conn.prepare_single("SELECT count(*) FROM users").unwrap();
2091        assert_eq!(verify_stmt.step(None).unwrap(), TursoStatusCode::Row);
2092        assert_eq!(verify_stmt.row_value(0).unwrap().as_int(), Some(1));
2093    }
2094
2095    #[cfg(clt_turso_feature = "encryption")]
2096    mod encryption_tests {
2097        use super::*;
2098        use tempfile::NamedTempFile;
2099
2100        const TEST_CIPHER: &str = "aes256gcm";
2101        const TEST_HEXKEY: &str =
2102            "000102030405060708090a0b0c0d0e0f101112131415161718191a1b1c1d1e1f";
2103        const WRONG_HEXKEY: &str =
2104            "ffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff";
2105
2106        fn create_encryption_opts() -> crate::turso_sdk_kit::rsapi::EncryptionOpts {
2107            crate::turso_sdk_kit::rsapi::EncryptionOpts {
2108                cipher: TEST_CIPHER.to_string(),
2109                hexkey: TEST_HEXKEY.to_string(),
2110            }
2111        }
2112
2113        fn assert_integer(value: turso_core::Value, expected: i64) {
2114            match value {
2115                turso_core::Value::Numeric(turso_core::Numeric::Integer(i)) => {
2116                    assert_eq!(i, expected)
2117                }
2118                _ => panic!("Expected integer {expected}, got {value:?}"),
2119            }
2120        }
2121
2122        #[test]
2123        fn test_encryption() {
2124            let temp_file = NamedTempFile::new().unwrap();
2125            let db_path = temp_file.path().to_str().unwrap();
2126
2127            // 1. Create encrypted database and insert data
2128            {
2129                let db = TursoDatabase::new(TursoDatabaseConfig {
2130                    path: db_path.to_string(),
2131                    experimental_features: Some("encryption".to_string()),
2132                    async_io: false,
2133                    encryption: Some(create_encryption_opts()),
2134                    vfs: None,
2135                    io: None,
2136                    db_file: None,
2137                });
2138                let result = db.open().unwrap();
2139                assert!(!result.is_io());
2140                let conn = db.connect().unwrap();
2141
2142                let mut stmt = conn
2143                    .prepare_single("CREATE TABLE test (id INTEGER PRIMARY KEY, value TEXT)")
2144                    .unwrap();
2145                stmt.execute(None).unwrap();
2146
2147                let mut stmt = conn
2148                    .prepare_single("INSERT INTO test (id, value) VALUES (1, 'secret_data')")
2149                    .unwrap();
2150                stmt.execute(None).unwrap();
2151
2152                // Checkpoint to ensure data is written to main db file
2153                let mut stmt = conn
2154                    .prepare_single("PRAGMA wal_checkpoint(TRUNCATE)")
2155                    .unwrap();
2156                stmt.execute(None).unwrap();
2157            }
2158
2159            // 2. Verify data is encrypted on disk
2160            let content = std::fs::read(db_path).unwrap();
2161            assert!(content.len() > 1024);
2162            assert!(
2163                !content.windows(11).any(|w| w == b"secret_data"),
2164                "Plaintext should not appear in encrypted database file"
2165            );
2166
2167            // 3. Reopen with correct key and verify data
2168            {
2169                let db = TursoDatabase::new(TursoDatabaseConfig {
2170                    path: db_path.to_string(),
2171                    experimental_features: Some("encryption".to_string()),
2172                    async_io: false,
2173                    encryption: Some(create_encryption_opts()),
2174                    vfs: None,
2175                    io: None,
2176                    db_file: None,
2177                });
2178                let result = db.open().unwrap();
2179                assert!(!result.is_io());
2180                let conn = db.connect().unwrap();
2181
2182                let mut stmt = conn
2183                    .prepare_single("SELECT id, value FROM test WHERE id = 1")
2184                    .unwrap();
2185                assert_eq!(stmt.step(None).unwrap(), TursoStatusCode::Row);
2186                assert_integer(stmt.row_value(0).unwrap(), 1);
2187                assert_eq!(stmt.row_value(1).unwrap().to_text(), Some("secret_data"));
2188            }
2189
2190            // 4. Verify opening with wrong key fails
2191            {
2192                let db = TursoDatabase::new(TursoDatabaseConfig {
2193                    path: db_path.to_string(),
2194                    experimental_features: Some("encryption".to_string()),
2195                    async_io: false,
2196                    encryption: Some(crate::turso_sdk_kit::rsapi::EncryptionOpts {
2197                        cipher: TEST_CIPHER.to_string(),
2198                        hexkey: WRONG_HEXKEY.to_string(),
2199                    }),
2200                    vfs: None,
2201                    io: None,
2202                    db_file: None,
2203                });
2204                assert!(db.open().is_err(), "Opening with wrong key should fail");
2205            }
2206
2207            // 5. Verify opening without encryption fails
2208            {
2209                let db = TursoDatabase::new(TursoDatabaseConfig {
2210                    path: db_path.to_string(),
2211                    experimental_features: Some("encryption".to_string()),
2212                    async_io: false,
2213                    encryption: None,
2214                    vfs: None,
2215                    io: None,
2216                    db_file: None,
2217                });
2218                let result = db.open();
2219                println!("result: {result:?}");
2220                assert!(
2221                    result.is_err(),
2222                    "Opening encrypted database without key should fail"
2223                );
2224            }
2225        }
2226    }
2227
2228    /// Reproducer: stale DATABASE_MANAGER entry when old TursoDatabase/TursoConnection
2229    /// haven't been GC'd (dropped) before reopening at the same path.
2230    ///
2231    /// Steps (mirrors the React Native bug report):
2232    ///   1. Open database A via SDK, create table "cache", close connection
2233    ///   2. Copy A.db → B.db, delete A.db
2234    ///   3. Open a *new* database at path A.db — while old db_a/conn_a still alive
2235    ///   4. CREATE TABLE cache should succeed (A.db is fresh) but fails with
2236    ///      "table cache already exists" because the registry returned the stale Database
2237    #[test]
2238    pub fn test_stale_registry_with_live_sdk_handles() {
2239        let tmp_dir = tempfile::TempDir::new().unwrap();
2240        let path_a = tmp_dir.path().join("A.db");
2241        let path_b = tmp_dir.path().join("B.db");
2242
2243        // 1. Open database A via SDK and create a table.
2244        let db_a = TursoDatabase::new(TursoDatabaseConfig {
2245            path: path_a.to_str().unwrap().to_string(),
2246            experimental_features: None,
2247            async_io: false,
2248            encryption: None,
2249            vfs: None,
2250            io: None,
2251            db_file: None,
2252        });
2253        let _ = db_a.open().unwrap();
2254        let conn_a = db_a.connect().unwrap();
2255
2256        let mut stmt = conn_a
2257            .prepare_single("CREATE TABLE cache(x INTEGER)")
2258            .unwrap();
2259        assert_eq!(stmt.execute(None).unwrap().status, TursoStatusCode::Done);
2260        drop(stmt);
2261
2262        // Close the connection but do NOT drop conn_a or db_a — simulates
2263        // the JS GC not having collected them yet.
2264        conn_a.close().unwrap();
2265
2266        // 2. Copy A.db → B.db, then delete A.db (and WAL/SHM files).
2267        std::fs::copy(&path_a, &path_b).unwrap();
2268        std::fs::remove_file(&path_a).unwrap();
2269        for ext in &["-wal", "-shm"] {
2270            let src = tmp_dir.path().join(format!("A.db{ext}"));
2271            let dst = tmp_dir.path().join(format!("B.db{ext}"));
2272            if src.exists() {
2273                std::fs::copy(&src, &dst).unwrap();
2274                std::fs::remove_file(&src).unwrap();
2275            }
2276        }
2277
2278        // 3. Open a new database at the same path A.db.
2279        //    The old db_a and conn_a are still alive — this is the key difference
2280        //    from test_sdk_close_finalizes_leaked_statements which drops everything.
2281        let db_a2 = TursoDatabase::new(TursoDatabaseConfig {
2282            path: path_a.to_str().unwrap().to_string(),
2283            experimental_features: None,
2284            async_io: false,
2285            encryption: None,
2286            vfs: None,
2287            io: None,
2288            db_file: None,
2289        });
2290        let _ = db_a2.open().unwrap();
2291        let conn_a2 = db_a2.connect().unwrap();
2292
2293        // 4. A.db should be a fresh empty database — CREATE TABLE cache must succeed.
2294        let mut stmt2 = conn_a2
2295            .prepare_single("CREATE TABLE cache(x INTEGER)")
2296            .expect("prepare should succeed on fresh database");
2297        let result = stmt2.execute(None);
2298        assert_eq!(
2299            result.unwrap().status,
2300            TursoStatusCode::Done,
2301            "CREATE TABLE cache on a fresh A.db should succeed — \
2302             stale DATABASE_MANAGER entry returned the old Database"
2303        );
2304
2305        // Cleanup: drop old handles (simulates eventual GC).
2306        drop(conn_a);
2307        drop(db_a);
2308    }
2309
2310    /// Regression test: connection.close() must finalize all outstanding statements
2311    /// to break the Statement → Arc<Connection> → Arc<Database> chain that keeps the
2312    /// database alive in DATABASE_MANAGER after a file rename.
2313    #[test]
2314    pub fn test_close_finalizes_outstanding_statements() {
2315        let db = TursoDatabase::new(TursoDatabaseConfig {
2316            path: ":memory:".to_string(),
2317            experimental_features: None,
2318            async_io: false,
2319            encryption: None,
2320            vfs: None,
2321            io: None,
2322            db_file: None,
2323        });
2324        let result = db.open().unwrap();
2325        assert!(!result.is_io());
2326
2327        let conn = db.connect().unwrap();
2328
2329        // Create a statement but do NOT finalize or drop it
2330        let mut stmt = conn.prepare_single("SELECT 1").unwrap();
2331        assert_eq!(stmt.step(None).unwrap(), TursoStatusCode::Row);
2332
2333        // close() should finalize the outstanding statement
2334        conn.close().unwrap();
2335
2336        // The statement should now be finalized — using it returns an error
2337        let result = stmt.step(None);
2338        assert!(result.is_err());
2339        match result.unwrap_err() {
2340            TursoError::Misuse(msg) => assert_eq!(msg, FINALIZED_ERR),
2341            other => panic!("expected Misuse error, got: {other:?}"),
2342        }
2343    }
2344
2345    /// Test that finalize() sets the statement handle to None, making subsequent
2346    /// operations return "statement has been finalized".
2347    #[test]
2348    pub fn test_finalize_disposes_statement() {
2349        let db = TursoDatabase::new(TursoDatabaseConfig {
2350            path: ":memory:".to_string(),
2351            experimental_features: None,
2352            async_io: false,
2353            encryption: None,
2354            vfs: None,
2355            io: None,
2356            db_file: None,
2357        });
2358        let result = db.open().unwrap();
2359        assert!(!result.is_io());
2360
2361        let conn = db.connect().unwrap();
2362        let mut stmt = conn.prepare_single("SELECT 1").unwrap();
2363
2364        // Finalize the statement
2365        assert_eq!(stmt.finalize(None).unwrap(), TursoStatusCode::Done);
2366
2367        // All operations should now return "statement has been finalized"
2368        assert!(stmt.step(None).is_err());
2369        assert!(stmt.execute(None).is_err());
2370        assert!(stmt.reset().is_err());
2371        assert!(stmt.run_io().is_err());
2372        assert!(stmt.bind_positional(1, Value::Null).is_err());
2373        assert_eq!(stmt.n_change(), 0);
2374        assert_eq!(stmt.column_count(), 0);
2375        assert_eq!(stmt.parameters_count(), 0);
2376    }
2377}