Skip to main content

turso_orm_driver/
database.rs

1//! The pooled database handle, modeled by [`Database`].
2//!
3//! Turso connections are cheap to open but not free, and the engine
4//! serialises work on a single connection, so the handle keeps a bounded
5//! pool of them. The same pool serves the serverless client, where a
6//! connection is a server-side session that carries transaction state
7//! between HTTP requests. The pool is a semaphore for the slot count plus an idle
8//! list: a borrower takes a permit, pops an idle connection or opens a new
9//! one with `db.connect()`, and the connection returns to the idle list on
10//! drop. Slots are never multiplied by cloning a connection,
11//! because clones share one engine connection and would serialise on it.
12//!
13//! A connection goes back to the idle list only when it is in autocommit
14//! mode. Anything else — a transaction dropped without commit or rollback,
15//! a statement left mid-way — is a state the next borrower must not
16//! inherit, so the connection is dropped instead and the engine rolls back
17//! whatever it held; a fresh one is opened on the next acquire.
18//!
19//! This module owns the engine handle, the pool and the per-connection
20//! setup. Statement execution lives in `crate::executor`, the connection
21//! traits in `crate::connection` and transactions in `crate::transaction`.
22//!
23//! - [`Database`]: the cloneable handle;
24//! - [`PooledConnection`]: a checked-out connection that returns itself on
25//!   drop;
26//! - [`retry_busy`]: the backoff loop used where a busy error is expected.
27
28use std::fmt;
29use std::future::Future;
30use std::ops::Deref;
31use std::sync::atomic::{AtomicBool, Ordering};
32use std::sync::{Arc, Mutex};
33use std::time::Duration;
34
35use tokio::sync::{OwnedSemaphorePermit, Semaphore};
36
37use crate::error::{Error, Result};
38use crate::executor::Conn;
39use crate::options::{ConnectOptions, Source};
40use turso_sql::Statement;
41
42/// The engine handle behind the pool.
43enum Engine {
44    /// A local file or in-memory database.
45    Local(turso::Database),
46    /// An embedded replica synchronised with Turso Cloud.
47    #[cfg(feature = "sync")]
48    Sync(turso::sync::Database),
49    /// A Turso Cloud database reached over HTTP.
50    #[cfg(feature = "serverless")]
51    Remote(turso_serverless::Database),
52}
53
54/// The state shared by every clone of a [`Database`].
55pub(crate) struct Inner {
56    /// The engine handle.
57    engine: Engine,
58    /// The options the database was opened with.
59    pub(crate) options: ConnectOptions,
60    /// Connections that are open and not checked out.
61    idle: Mutex<Vec<Conn>>,
62    /// One permit per pool slot, whether the slot's connection is idle or
63    /// not yet opened.
64    permits: Arc<Semaphore>,
65}
66
67/// A Turso database with a pool of connections.
68///
69/// Cloning is cheap; every clone shares the same engine handle and pool.
70#[derive(Clone)]
71pub struct Database {
72    /// The shared state.
73    pub(crate) inner: Arc<Inner>,
74}
75
76impl fmt::Debug for Database {
77    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
78        f.debug_struct("Database")
79            .field("source", &self.inner.options.source)
80            .field("max_connections", &self.inner.options.max_connections)
81            .finish_non_exhaustive()
82    }
83}
84
85impl Database {
86    /// Opens the database described by `options`.
87    ///
88    /// One connection is opened eagerly and dropped back into the pool so
89    /// that a bad path, a wrong key or a rejected pragma fails here rather
90    /// than on the first query.
91    ///
92    /// # Errors
93    ///
94    /// Returns [`Error::InvalidOptions`] when encryption is combined with a
95    /// sync source, and [`Error::Turso`] with [`ErrorKind::Connection`] or
96    /// [`ErrorKind::Other`] when the engine cannot open the database, open
97    /// the first connection or apply its pragmas.
98    ///
99    /// [`ErrorKind::Connection`]: crate::ErrorKind::Connection
100    /// [`ErrorKind::Other`]: crate::ErrorKind::Other
101    pub async fn connect(options: impl Into<ConnectOptions>) -> Result<Self> {
102        let options = options.into();
103        let engine = match &options.source {
104            Source::Memory => Engine::Local(options.local_builder(":memory:").build().await?),
105            Source::File(path) => Engine::Local(
106                options
107                    .local_builder(&path.to_string_lossy())
108                    .build()
109                    .await?,
110            ),
111            #[cfg(feature = "sync")]
112            Source::Sync(sync) => {
113                // The sync builder has no encryption hook, so the option
114                // would be silently ignored; refuse it instead.
115                if options.encryption.is_some() {
116                    return Err(Error::InvalidOptions(
117                        "local encryption is not supported together with sync".into(),
118                    ));
119                }
120                let mut builder = turso::sync::Builder::new_remote(&sync.path.to_string_lossy())
121                    .with_remote_url(&sync.remote_url)
122                    .bootstrap_if_empty(sync.bootstrap_if_empty);
123                if let Some(token) = &sync.auth_token {
124                    builder = builder.with_auth_token(token);
125                }
126                Engine::Sync(builder.build().await?)
127            }
128            #[cfg(feature = "serverless")]
129            Source::Remote(remote) => {
130                if options.encryption.is_some() {
131                    return Err(Error::InvalidOptions(
132                        "local encryption does not apply to a remote database".into(),
133                    ));
134                }
135                let mut builder = turso_serverless::Builder::new_remote(remote.url.clone());
136                if let Some(token) = &remote.auth_token {
137                    builder = builder.with_auth_token(token.clone());
138                }
139                if let Some(key) = &remote.remote_encryption_key {
140                    builder = builder.with_remote_encryption_key(key.clone());
141                }
142                Engine::Remote(builder.build().await?)
143            }
144        };
145        let db = Self {
146            inner: Arc::new(Inner {
147                engine,
148                permits: Arc::new(Semaphore::new(options.max_connections)),
149                idle: Mutex::new(Vec::with_capacity(options.max_connections)),
150                options,
151            }),
152        };
153        // Warm one connection so that misconfiguration fails fast; dropping
154        // it parks it in the idle list for the first real borrower.
155        drop(db.acquire().await?);
156        Ok(db)
157    }
158
159    /// The options this database was opened with.
160    pub fn options(&self) -> &ConnectOptions {
161        &self.inner.options
162    }
163
164    /// Checks that the database answers queries.
165    ///
166    /// # Errors
167    ///
168    /// Returns [`Error::PoolTimeout`] when no connection is free within the
169    /// acquire timeout, [`Error::Misuse`] when the pool is closed, and
170    /// [`Error::Turso`] when a connection cannot be opened or the probe
171    /// query fails.
172    pub async fn ping(&self) -> Result<()> {
173        let conn = self.acquire().await?;
174        crate::executor::query_one(&conn, &Statement::from_string("SELECT 1")).await?;
175        Ok(())
176    }
177
178    /// Pushes local changes to Turso Cloud — embedded replica only.
179    ///
180    /// # Errors
181    ///
182    /// Returns [`Error::InvalidOptions`] when this database is not a
183    /// replica, and [`Error::Turso`] when the sync fails.
184    #[cfg(feature = "sync")]
185    #[cfg_attr(docsrs, doc(cfg(feature = "sync")))]
186    pub async fn push(&self) -> Result<()> {
187        match &self.inner.engine {
188            Engine::Sync(db) => Ok(db.push().await?),
189            _ => Err(Error::InvalidOptions("not an embedded replica".into())),
190        }
191    }
192
193    /// Pulls remote changes from Turso Cloud — embedded replica only — and
194    /// returns whether anything changed.
195    ///
196    /// # Errors
197    ///
198    /// Returns [`Error::InvalidOptions`] when this database is not a
199    /// replica, and [`Error::Turso`] when the sync fails.
200    #[cfg(feature = "sync")]
201    #[cfg_attr(docsrs, doc(cfg(feature = "sync")))]
202    pub async fn pull(&self) -> Result<bool> {
203        match &self.inner.engine {
204            Engine::Sync(db) => Ok(db.pull().await?),
205            _ => Err(Error::InvalidOptions("not an embedded replica".into())),
206        }
207    }
208
209    /// Takes a connection from the pool, opening a new one if none is idle.
210    ///
211    /// The permit is held by the returned guard, so the pool never has more
212    /// than `max_connections` connections checked out or idle.
213    ///
214    /// # Errors
215    ///
216    /// Returns [`Error::PoolTimeout`] when no permit is free within the
217    /// acquire timeout, [`Error::Misuse`] when the semaphore is closed or
218    /// the idle-list mutex is poisoned, and [`Error::Turso`] when a new
219    /// connection cannot be opened or configured.
220    pub(crate) async fn acquire(&self) -> Result<PooledConnection> {
221        let permit = tokio::time::timeout(
222            self.inner.options.acquire_timeout,
223            Arc::clone(&self.inner.permits).acquire_owned(),
224        )
225        .await
226        .map_err(|_| Error::PoolTimeout)?
227        .map_err(|_| Error::Misuse("connection pool closed".into()))?;
228
229        // The lock is released before the await below so that opening a
230        // connection never blocks other borrowers from popping idle ones.
231        let idle = self
232            .inner
233            .idle
234            .lock()
235            .map_err(|_| Error::Misuse("pool mutex poisoned".into()))?
236            .pop();
237        let conn = match idle {
238            Some(conn) => conn,
239            None => self.open_connection().await?,
240        };
241        Ok(PooledConnection {
242            conn: Some(conn),
243            pool: Arc::clone(&self.inner),
244            _permit: permit,
245            discard: AtomicBool::new(false),
246        })
247    }
248
249    /// Opens a new engine connection and applies the per-connection
250    /// settings.
251    ///
252    /// Pragmas are per connection in SQLite, so every connection the pool
253    /// opens must be configured the same way or borrowers would observe
254    /// different behaviour depending on which slot they get.
255    ///
256    /// # Errors
257    ///
258    /// Returns [`Error::Turso`] when the engine cannot open the connection
259    /// or rejects one of the settings.
260    async fn open_connection(&self) -> Result<Conn> {
261        let conn = match &self.inner.engine {
262            Engine::Local(db) => Conn::Embedded(db.connect()?),
263            #[cfg(feature = "sync")]
264            Engine::Sync(db) => Conn::Embedded(db.connect().await?),
265            #[cfg(feature = "serverless")]
266            Engine::Remote(db) => Conn::Remote(db.connect()?),
267        };
268        let options = &self.inner.options;
269        if let Some(timeout) = options.busy_timeout {
270            conn.busy_timeout(timeout)?;
271        }
272        if options.foreign_keys {
273            conn.pragma_update("foreign_keys", "ON").await?;
274        }
275        // MVCC is a property of the local engine; a remote session has no
276        // journal to switch.
277        if options.mvcc && matches!(conn, Conn::Embedded(_)) {
278            conn.pragma_update("journal_mode", "'mvcc'").await?;
279        }
280        for (name, value) in &options.pragmas {
281            conn.pragma_update(name, value).await?;
282        }
283        Ok(conn)
284    }
285}
286
287/// A connection checked out of the pool.
288///
289/// The connection returns to the idle list on drop unless it was discarded
290/// or is no longer in autocommit mode; the permit is released either way.
291pub(crate) struct PooledConnection {
292    /// The connection; `None` only once `drop` has taken it.
293    conn: Option<Conn>,
294    /// The pool to return the connection to.
295    pool: Arc<Inner>,
296    /// The slot permit, released when the guard is dropped.
297    _permit: OwnedSemaphorePermit,
298    /// Whether the connection must be dropped instead of returned.
299    discard: AtomicBool,
300}
301
302impl PooledConnection {
303    /// Marks the connection so that it is not returned to the pool.
304    ///
305    /// Used when the connection's state is unknown, for example because a
306    /// transaction was dropped without commit or rollback.
307    pub(crate) fn discard(&self) {
308        self.discard.store(true, Ordering::Release);
309    }
310
311    /// The options of the pool this connection belongs to.
312    pub(crate) fn options(&self) -> &ConnectOptions {
313        &self.pool.options
314    }
315}
316
317impl Deref for PooledConnection {
318    type Target = Conn;
319
320    /// The underlying engine connection.
321    ///
322    /// # Panics
323    ///
324    /// When called after `drop` has taken the connection, which safe code
325    /// cannot do since the field is only emptied inside `Drop`.
326    fn deref(&self) -> &Self::Target {
327        self.conn.as_ref().expect("connection present until drop")
328    }
329}
330
331impl Drop for PooledConnection {
332    fn drop(&mut self) {
333        // A connection that is not in autocommit mode still holds a
334        // transaction — typically one that was dropped without commit or
335        // rollback — and must not be handed to the next borrower; dropping
336        // it makes the engine roll back. An error from `is_autocommit` is
337        // treated the same way, since the connection's state is unknown.
338        if let Some(conn) = self.conn.take()
339            && !self.discard.load(Ordering::Acquire)
340            && conn.is_autocommit().unwrap_or(false)
341            && let Ok(mut idle) = self.pool.idle.lock()
342        {
343            idle.push(conn);
344        }
345    }
346}
347
348impl fmt::Debug for PooledConnection {
349    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
350        f.debug_struct("PooledConnection").finish_non_exhaustive()
351    }
352}
353
354/// Retries `op` with exponential backoff while Turso reports lock
355/// contention, for at most `budget`.
356///
357/// The engine's own busy handler only covers lock waits inside a statement;
358/// `BEGIN IMMEDIATE` under contention and MVCC commit conflicts surface as
359/// busy errors immediately, so callers that expect them wrap the call in
360/// this loop.
361///
362/// # Errors
363///
364/// Returns whatever `op` returns once it fails with a non-busy error or the
365/// budget is exhausted, so [`Error::Turso`] with
366/// [`ErrorKind::Busy`](crate::ErrorKind::Busy) is what a persistent lock
367/// produces.
368pub(crate) async fn retry_busy<T, F, Fut>(budget: Duration, mut op: F) -> Result<T>
369where
370    F: FnMut() -> Fut,
371    Fut: Future<Output = Result<T>>,
372{
373    let start = std::time::Instant::now();
374    // The delay starts small enough not to add latency to a lock that is
375    // about to clear and doubles up to a cap that keeps the loop responsive
376    // once the budget runs into seconds.
377    let mut delay = Duration::from_millis(5);
378    loop {
379        match op().await {
380            Err(err) if err.is_busy() && start.elapsed() < budget => {
381                tracing::debug!(%err, ?delay, "busy, retrying");
382                tokio::time::sleep(delay).await;
383                delay = (delay * 2).min(Duration::from_millis(250));
384            }
385            other => return other,
386        }
387    }
388}