Skip to main content

turso_orm_driver/
executor.rs

1//! Statement execution on a raw engine connection, shared by the pool and transactions.
2//!
3//! Both [`Database`] and [`Transaction`] end up with a [`Conn`] and a
4//! [`Statement`]; everything from there — preparing through the engine's
5//! statement cache, binding the values, collecting or streaming rows and
6//! reading the change counters — is identical and lives here once. The
7//! module owns the [`Row`] type as well, because the column index a row
8//! carries is built from the engine's result set at execution time.
9//!
10//! Two engines can sit behind a [`Conn`]: the embedded `turso` client, for
11//! in-memory and file databases and embedded replicas, and the
12//! `turso_serverless` HTTP client for Turso Cloud, behind the `serverless`
13//! feature. The two crates expose the same method names on purpose, so the
14//! engine-specific primitives are written once as a macro and instantiated
15//! per crate; [`Conn`] dispatches to the right instance. Rows carry
16//! `turso_sql::Value` rather than either engine's value type, which is what
17//! keeps decoding independent of the engine.
18//!
19//! Statements are prepared with `prepare_cached`, so the same SQL text
20//! reuses a compiled statement on the same connection; that is why the SQL
21//! layer binds paging values as parameters instead of inlining them. A
22//! cached statement is shared, so a query that stops early drains the
23//! remaining rows rather than leaving a cursor open on it.
24//!
25//! Decoding is not done here: a [`Row`] keeps the storage values and
26//! [`FromValue`] converts on access, so the same column can be read as
27//! different Rust types.
28//!
29//! - [`Conn`]: the engine connection behind the pool;
30//! - [`Row`] and [`ExecResult`]: what callers get back;
31//! - [`RowStream`]: the boxed stream type of [`StreamTrait`];
32//! - the `pub(crate)` functions: the execution primitives.
33//!
34//! [`Database`]: crate::Database
35//! [`Transaction`]: crate::Transaction
36//! [`StreamTrait`]: crate::StreamTrait
37
38use std::collections::HashMap;
39use std::fmt;
40use std::pin::Pin;
41use std::sync::Arc;
42use std::time::Duration;
43
44use futures_util::Stream;
45use turso_sql::{Statement, Value};
46
47use crate::decode::FromValue;
48use crate::error::{Error, Result};
49
50/// The result of a statement that does not return rows.
51#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
52pub struct ExecResult {
53    /// The value of `last_insert_rowid()` after the statement.
54    pub last_insert_id: i64,
55    /// The number of rows changed.
56    pub rows_affected: u64,
57}
58
59/// A result row.
60///
61/// Values are kept in their storage class and decoded on access with
62/// [`FromValue`], so the same column can be read as `i64`, `bool` or
63/// `String`. The column index is shared between all rows of one result set.
64#[derive(Clone)]
65pub struct Row {
66    /// The column names and lookup index, shared across the result set.
67    columns: Arc<Columns>,
68    /// The values, in column order.
69    values: Vec<Value>,
70}
71
72/// The column names of a result set with a case-insensitive lookup index.
73struct Columns {
74    /// The names in `SELECT` order, as the engine reports them.
75    names: Vec<String>,
76    /// Lower-cased name to position.
77    index: HashMap<String, usize>,
78}
79
80impl Columns {
81    /// Builds the shared column index for a result set.
82    ///
83    /// Names are indexed lower-cased because SQL identifiers are case
84    /// insensitive and callers often write `"ID"` for a column declared as
85    /// `id`. The last occurrence of a duplicated name wins.
86    fn new(names: Vec<String>) -> Arc<Self> {
87        let index = names
88            .iter()
89            .enumerate()
90            .map(|(i, n)| (n.to_ascii_lowercase(), i))
91            .collect();
92        Arc::new(Self { names, index })
93    }
94}
95
96impl fmt::Debug for Row {
97    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
98        let mut m = f.debug_map();
99        for (name, value) in self.columns.names.iter().zip(&self.values) {
100            m.entry(name, value);
101        }
102        m.finish()
103    }
104}
105
106/// A column lookup: by name (`&str`) or by position (`usize`).
107pub trait ColumnIndex: fmt::Display + Copy {
108    /// Resolves to a position in the row, or `None` when absent.
109    fn resolve(self, row: &Row) -> Option<usize>;
110}
111
112impl ColumnIndex for usize {
113    fn resolve(self, row: &Row) -> Option<usize> {
114        (self < row.values.len()).then_some(self)
115    }
116}
117
118impl ColumnIndex for &str {
119    /// Resolves the name as given first, then lower-cased, so an exact
120    /// match wins when a result set has names differing only in case.
121    fn resolve(self, row: &Row) -> Option<usize> {
122        row.columns
123            .index
124            .get(self)
125            .or_else(|| row.columns.index.get(&self.to_ascii_lowercase()))
126            .copied()
127    }
128}
129
130impl Row {
131    /// The column names in `SELECT` order.
132    pub fn columns(&self) -> &[String] {
133        &self.columns.names
134    }
135
136    /// The number of columns.
137    pub fn len(&self) -> usize {
138        self.values.len()
139    }
140
141    /// Whether the row has no columns.
142    pub fn is_empty(&self) -> bool {
143        self.values.is_empty()
144    }
145
146    /// Whether a column exists.
147    pub fn has(&self, column: &str) -> bool {
148        column.resolve(self).is_some()
149    }
150
151    /// Decodes a column.
152    ///
153    /// # Errors
154    ///
155    /// Returns [`Error::Decode`] when the column does not exist or the value
156    /// cannot be converted to `T`.
157    pub fn get<T: FromValue>(&self, column: impl ColumnIndex) -> Result<T> {
158        let idx = column
159            .resolve(self)
160            .ok_or_else(|| Error::decode(column, T::TYPE_NAME, "no such column"))?;
161        T::from_value(self.values[idx].clone(), &self.columns.names[idx])
162    }
163
164    /// Decodes a column, returning `None` when it does not exist.
165    ///
166    /// # Errors
167    ///
168    /// Returns [`Error::Decode`] when the value cannot be converted to `T`.
169    pub fn try_get<T: FromValue>(&self, column: impl ColumnIndex) -> Result<Option<T>> {
170        match column.resolve(self) {
171            None => Ok(None),
172            Some(idx) => {
173                T::from_value(self.values[idx].clone(), &self.columns.names[idx]).map(Some)
174            }
175        }
176    }
177
178    /// The raw storage value of a column, or `None` when absent.
179    pub fn raw(&self, column: impl ColumnIndex) -> Option<&Value> {
180        column.resolve(self).map(|i| &self.values[i])
181    }
182
183    /// Iterates over `(name, value)` pairs in column order.
184    pub fn iter(&self) -> impl Iterator<Item = (&str, &Value)> {
185        self.columns
186            .names
187            .iter()
188            .map(String::as_str)
189            .zip(self.values.iter())
190    }
191}
192
193/// A boxed, sendable stream of rows.
194pub type RowStream<'a> = Pin<Box<dyn Stream<Item = Result<Row>> + Send + 'a>>;
195
196/// Generates the execution primitives for one engine crate.
197///
198/// The embedded and the serverless clients share their method names and
199/// value shapes, so one body serves both; only the crate path differs.
200/// Each instance is a private module holding the conversions between
201/// `turso_sql::Value` and the engine's value type and the primitives
202/// [`Conn`] dispatches to.
203macro_rules! engine_module {
204    ($(#[$meta:meta])* $name:ident, $engine:ident) => {
205        $(#[$meta])*
206        mod $name {
207            use std::sync::Arc;
208
209            use $engine::params_from_iter;
210            use turso_sql::{Statement, Value};
211
212            use super::{Columns, ExecResult, Row, RowStream};
213            use crate::error::Result;
214
215            /// Converts an engine value into the SQL layer's value type.
216            fn from_engine(value: $engine::Value) -> Value {
217                match value {
218                    $engine::Value::Null => Value::Null,
219                    $engine::Value::Integer(n) => Value::Integer(n),
220                    $engine::Value::Real(f) => Value::Real(f),
221                    $engine::Value::Text(s) => Value::Text(s),
222                    $engine::Value::Blob(b) => Value::Blob(b),
223                }
224            }
225
226            /// Converts a SQL layer value into the engine's value type.
227            fn to_engine(value: Value) -> $engine::Value {
228                match value {
229                    Value::Null => $engine::Value::Null,
230                    Value::Integer(n) => $engine::Value::Integer(n),
231                    Value::Real(f) => $engine::Value::Real(f),
232                    Value::Text(s) => $engine::Value::Text(s),
233                    Value::Blob(b) => $engine::Value::Blob(b),
234                }
235            }
236
237            /// The statement's bound values as engine parameters, in
238            /// placeholder order.
239            fn params(statement: &Statement) -> Vec<$engine::Value> {
240                statement.values.iter().cloned().map(to_engine).collect()
241            }
242
243            /// Copies an engine row into an owned [`Row`] sharing `columns`.
244            fn row_from(columns: &Arc<Columns>, row: &$engine::Row) -> Result<Row> {
245                let mut values = Vec::with_capacity(columns.names.len());
246                for i in 0..columns.names.len() {
247                    values.push(from_engine(row.get_value(i)?));
248                }
249                Ok(Row {
250                    columns: Arc::clone(columns),
251                    values,
252                })
253            }
254
255            /// Prepares `sql` through the connection's statement cache.
256            async fn prepare(conn: &$engine::Connection, sql: &str) -> Result<$engine::Statement> {
257                Ok(conn.prepare_cached(sql).await?)
258            }
259
260            /// Runs a query and collects every row.
261            pub(super) async fn query_all(
262                conn: &$engine::Connection,
263                statement: &Statement,
264            ) -> Result<Vec<Row>> {
265                let mut stmt = prepare(conn, &statement.sql).await?;
266                let mut rows = stmt.query(params_from_iter(params(statement))).await?;
267                let columns = Columns::new(rows.column_names());
268                let mut out = Vec::new();
269                while let Some(row) = rows.next().await? {
270                    out.push(row_from(&columns, &row)?);
271                }
272                Ok(out)
273            }
274
275            /// Runs a query and returns its first row, if any.
276            pub(super) async fn query_one(
277                conn: &$engine::Connection,
278                statement: &Statement,
279            ) -> Result<Option<Row>> {
280                let mut stmt = prepare(conn, &statement.sql).await?;
281                let mut rows = stmt.query(params_from_iter(params(statement))).await?;
282                let columns = Columns::new(rows.column_names());
283                let first = rows.next().await?;
284                // The remaining rows are drained so that the cached
285                // statement is left fully stepped rather than holding an
286                // open cursor, and its implicit read transaction, until it
287                // is next reused.
288                while rows.next().await?.is_some() {}
289                first.map(|row| row_from(&columns, &row)).transpose()
290            }
291
292            /// Runs a statement that returns no rows and reads the change
293            /// counters.
294            pub(super) async fn execute(
295                conn: &$engine::Connection,
296                statement: &Statement,
297            ) -> Result<ExecResult> {
298                let mut stmt = prepare(conn, &statement.sql).await?;
299                let rows_affected = stmt.execute(params_from_iter(params(statement))).await?;
300                Ok(ExecResult {
301                    last_insert_id: conn.last_insert_rowid(),
302                    rows_affected,
303                })
304            }
305
306            /// Runs one or more `;`-separated statements without
307            /// parameters. The batch API reports no change count.
308            pub(super) async fn execute_unprepared(
309                conn: &$engine::Connection,
310                sql: &str,
311            ) -> Result<ExecResult> {
312                conn.execute_batch(sql).await?;
313                Ok(ExecResult {
314                    last_insert_id: conn.last_insert_rowid(),
315                    rows_affected: 0,
316                })
317            }
318
319            /// Runs a single parameterless statement, for transaction
320            /// control.
321            pub(super) async fn execute_raw(conn: &$engine::Connection, sql: &str) -> Result<()> {
322                conn.execute(sql, ()).await?;
323                Ok(())
324            }
325
326            /// Sets a pragma on the connection.
327            pub(super) async fn pragma_update(
328                conn: &$engine::Connection,
329                name: &str,
330                value: &str,
331            ) -> Result<()> {
332                conn.pragma_update(name, value).await?;
333                Ok(())
334            }
335
336            /// Runs a query and streams its rows lazily; `holder` travels
337            /// inside the stream and is dropped with it.
338            pub(super) async fn stream<'a, H: Send + 'a>(
339                conn: &$engine::Connection,
340                statement: &Statement,
341                holder: H,
342            ) -> Result<RowStream<'a>> {
343                let mut stmt = prepare(conn, &statement.sql).await?;
344                let rows = stmt.query(params_from_iter(params(statement))).await?;
345                let columns = Columns::new(rows.column_names());
346                Ok(Box::pin(futures_util::stream::unfold(
347                    (rows, columns, holder),
348                    |(mut rows, columns, holder)| async move {
349                        match rows.next().await {
350                            Ok(Some(row)) => {
351                                Some((row_from(&columns, &row), (rows, columns, holder)))
352                            }
353                            Ok(None) => None,
354                            Err(e) => Some((Err(e.into()), (rows, columns, holder))),
355                        }
356                    },
357                )))
358            }
359        }
360    };
361}
362
363engine_module!(embedded, turso);
364engine_module!(
365    #[cfg(feature = "serverless")]
366    remote,
367    turso_serverless
368);
369
370/// An engine connection: embedded, or an HTTP session to Turso Cloud.
371///
372/// Both variants are cheap to clone, and clones share one engine
373/// connection, which is why the pool never multiplies slots by cloning.
374#[derive(Clone)]
375pub(crate) enum Conn {
376    /// A connection of the embedded engine: in-memory, file or replica.
377    Embedded(turso::Connection),
378    /// A session of the serverless client, one HTTP request per statement.
379    #[cfg(feature = "serverless")]
380    Remote(turso_serverless::Connection),
381}
382
383/// Dispatches one primitive call to the engine behind a [`Conn`].
384macro_rules! dispatch {
385    ($conn:expr, |$c:ident| $call:expr) => {
386        match $conn {
387            Conn::Embedded($c) => {
388                use embedded as engine;
389                $call
390            }
391            #[cfg(feature = "serverless")]
392            Conn::Remote($c) => {
393                use remote as engine;
394                $call
395            }
396        }
397    };
398}
399
400impl Conn {
401    /// Whether no explicit transaction is open on this connection.
402    ///
403    /// # Errors
404    ///
405    /// Returns the engine error when the state cannot be read.
406    pub(crate) fn is_autocommit(&self) -> Result<bool> {
407        match self {
408            Conn::Embedded(c) => Ok(c.is_autocommit()?),
409            #[cfg(feature = "serverless")]
410            Conn::Remote(c) => Ok(c.is_autocommit()?),
411        }
412    }
413
414    /// Sets the engine's own lock wait — embedded engine only, since an
415    /// HTTP session has no local lock to wait on.
416    ///
417    /// # Errors
418    ///
419    /// Returns the engine error when the timeout is rejected.
420    pub(crate) fn busy_timeout(&self, timeout: Duration) -> Result<()> {
421        match self {
422            Conn::Embedded(c) => Ok(c.busy_timeout(timeout)?),
423            #[cfg(feature = "serverless")]
424            Conn::Remote(_) => Ok(()),
425        }
426    }
427
428    /// Sets a pragma on the connection.
429    ///
430    /// # Errors
431    ///
432    /// Returns the engine error when the pragma is rejected.
433    pub(crate) async fn pragma_update(&self, name: &str, value: &str) -> Result<()> {
434        dispatch!(self, |c| engine::pragma_update(c, name, value).await)
435    }
436
437    /// Runs a single parameterless statement, for transaction control.
438    ///
439    /// # Errors
440    ///
441    /// Returns the engine error when the statement fails.
442    pub(crate) async fn execute_raw(&self, sql: &str) -> Result<()> {
443        dispatch!(self, |c| engine::execute_raw(c, sql).await)
444    }
445}
446
447/// Runs a query and collects every row.
448///
449/// # Errors
450///
451/// Returns the engine error when the statement cannot be prepared, bound or
452/// stepped.
453pub(crate) async fn query_all(conn: &Conn, statement: &Statement) -> Result<Vec<Row>> {
454    tracing::debug!(sql = %statement.sql, "query_all");
455    dispatch!(conn, |c| engine::query_all(c, statement).await)
456}
457
458/// Runs a query and returns its first row, if any.
459///
460/// # Errors
461///
462/// Returns the engine error when the statement cannot be prepared, bound or
463/// stepped.
464pub(crate) async fn query_one(conn: &Conn, statement: &Statement) -> Result<Option<Row>> {
465    tracing::debug!(sql = %statement.sql, "query_one");
466    dispatch!(conn, |c| engine::query_one(c, statement).await)
467}
468
469/// Runs a statement that returns no rows and reads the change counters.
470///
471/// # Errors
472///
473/// Returns the engine error when the statement cannot be prepared, bound or
474/// executed — with [`ErrorKind::Constraint`](crate::ErrorKind::Constraint)
475/// on a constraint violation.
476pub(crate) async fn execute(conn: &Conn, statement: &Statement) -> Result<ExecResult> {
477    tracing::debug!(sql = %statement.sql, "execute");
478    dispatch!(conn, |c| engine::execute(c, statement).await)
479}
480
481/// Runs one or more `;`-separated statements without parameters.
482///
483/// The engines' batch API does not report a change count, so
484/// `rows_affected` is always `0` here; `last_insert_id` is still read after
485/// the batch.
486///
487/// # Errors
488///
489/// Returns the engine error when any statement of the batch fails.
490pub(crate) async fn execute_unprepared(conn: &Conn, sql: &str) -> Result<ExecResult> {
491    tracing::debug!(sql, "execute_unprepared");
492    dispatch!(conn, |c| engine::execute_unprepared(c, sql).await)
493}
494
495/// Runs a query and streams its rows lazily.
496///
497/// `holder` is any value that must stay alive while the stream is consumed,
498/// for example the pooled connection; it is moved into the stream's state
499/// and dropped with it.
500///
501/// # Errors
502///
503/// Returns the engine error when the statement cannot be prepared or
504/// started. Errors while stepping are yielded as stream items.
505pub(crate) async fn stream<'a, H: Send + 'a>(
506    conn: &Conn,
507    statement: &Statement,
508    holder: H,
509) -> Result<RowStream<'a>> {
510    tracing::debug!(sql = %statement.sql, "stream");
511    dispatch!(conn, |c| engine::stream(c, statement, holder).await)
512}