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}