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}