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 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 pub path: String,
156
157 pub experimental_features: Option<String>,
160
161 pub async_io: bool,
164
165 pub encryption: Option<EncryptionOpts>,
168
169 pub vfs: Option<String>,
175
176 pub io: Option<Arc<dyn IO>>,
178
179 pub db_file: Option<Arc<dyn DatabaseStorage>>,
182}
183
184impl TursoDatabaseConfig {
185 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 _ => 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
231pub 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
248pub 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
263pub 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
280pub 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
292pub 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 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#[derive(Default, Clone, Copy)]
369pub enum TursoDatabaseOpenPhase {
370 #[default]
371 Init,
372 Opening,
373 Done,
374}
375
376pub 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 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 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 pub const fn version() -> &'static str {
640 "0.7.2-clt.1"
641 }
642 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 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 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 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 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 let io: Arc<dyn turso_core::IO> = self.open_vfs_io()?;
747
748 *self.io.lock().unwrap() = Some(io.clone());
750
751 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 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 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 let connection = db.connect_with_encryption(encryption_key)?;
849
850 Ok(TursoConnection::new(&self.config, connection))
851 }
852
853 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 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 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 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 pub fn set_busy_timeout(&self, duration: Duration) {
934 self.connection.set_busy_timeout(duration);
935 }
936 pub fn interrupt(&self) {
941 self.connection.interrupt();
942 }
943 pub fn set_query_timeout(&self, duration: Duration) {
946 self.connection.set_query_timeout(duration);
947 }
948 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 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 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 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 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 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 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 pub fn close(&self) -> Result<(), TursoError> {
1183 let mut stmts = self.stmts.lock().unwrap();
1188 for (_id, weak) in stmts.drain() {
1189 if let Some(handle) = weak.upgrade() {
1190 *handle.lock().unwrap() = None;
1193 }
1194 }
1195 self.connection.close()?;
1196 Ok(())
1197 }
1198
1199 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 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 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 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 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
1251pub(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
1259fn 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 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 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 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 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 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 #[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 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 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 #[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 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 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 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 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 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 *handle = None;
1534 Ok(TursoStatusCode::Done)
1535 }
1536 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 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 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 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 assert_eq!(
1601 config_with_features(None).database_opts(),
1602 turso_core::DatabaseOpts::new()
1603 );
1604
1605 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 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 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 {
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 let mut stmt = conn
2154 .prepare_single("PRAGMA wal_checkpoint(TRUNCATE)")
2155 .unwrap();
2156 stmt.execute(None).unwrap();
2157 }
2158
2159 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 {
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 {
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 {
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 #[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 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 conn_a.close().unwrap();
2265
2266 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 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 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 drop(conn_a);
2307 drop(db_a);
2308 }
2309
2310 #[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 let mut stmt = conn.prepare_single("SELECT 1").unwrap();
2331 assert_eq!(stmt.step(None).unwrap(), TursoStatusCode::Row);
2332
2333 conn.close().unwrap();
2335
2336 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]
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 assert_eq!(stmt.finalize(None).unwrap(), TursoStatusCode::Done);
2366
2367 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}