Skip to main content

teaql_provider_sqlite/
lib.rs

1use std::collections::{BTreeMap, BTreeSet};
2use std::future::Future;
3use std::pin::Pin;
4use std::str::FromStr;
5use std::sync::{Arc, Mutex, MutexGuard};
6
7use chrono::{DateTime, NaiveDate, NaiveDateTime};
8use rusqlite::types::{Value as SqliteValue, ValueRef};
9use rusqlite::{Connection, Row, params_from_iter};
10use rust_decimal::Decimal;
11use teaql_core::{
12    DataType, EntityDescriptor, Expr, InsertCommand, PropertyDescriptor, Record, SelectQuery,
13    UpdateCommand, Value,
14};
15use teaql_runtime::{
16    GraphNode, InternalIdGenerator, RawAuditEvent, RuntimeError, SchemaProvider, UserContext,
17};
18use teaql_sql::{
19    CompiledQuery, DatabaseKind, SqlCompileError, SqlDialect, SqlTransport,
20    quote_identifier_if_needed,
21};
22
23pub const DEFAULT_ID_SPACE_TABLE: &str = "teaql_id_space";
24
25#[derive(Debug, Default, Clone, Copy)]
26pub struct SqliteDialect;
27
28impl SqlDialect for SqliteDialect {
29    fn kind(&self) -> DatabaseKind {
30        DatabaseKind::Sqlite
31    }
32
33    fn quote_ident(&self, ident: &str) -> String {
34        quote_ident(ident)
35    }
36
37    fn placeholder(&self, _index: usize) -> String {
38        "?".to_owned()
39    }
40
41    fn schema_type_sql(
42        &self,
43        data_type: DataType,
44        property: &PropertyDescriptor,
45    ) -> Result<&'static str, SqlCompileError> {
46        match data_type {
47            DataType::Bool => Ok("BOOLEAN"),
48            DataType::I64 | DataType::U64 if property.is_id => Ok("INTEGER"),
49            DataType::I64 | DataType::U64 => Ok("INTEGER"),
50            DataType::F64 => Ok("REAL"),
51            DataType::Decimal => Ok("NUMERIC"),
52            DataType::Text => Ok("VARCHAR(255)"),
53            DataType::LargeText => Ok("TEXT"),
54            DataType::Json => Ok("JSON"),
55            DataType::Date => Ok("DATE"),
56            DataType::Timestamp => Ok("TIMESTAMP"),
57        }
58    }
59
60    fn compile_add_column(
61        &self,
62        entity: &EntityDescriptor,
63        property: &PropertyDescriptor,
64    ) -> Result<String, SqlCompileError> {
65        // SQLite does not support adding NOT NULL columns without a DEFAULT.
66        // Since TeaQL enforces nullability at the application layer, we can safely
67        // strip the NOT NULL constraint when adding columns to existing tables.
68        let def = self.column_definition_sql(property)?;
69        let def_without_not_null = def.replace(" NOT NULL", "");
70
71        Ok(format!(
72            "ALTER TABLE {} ADD COLUMN {}",
73            self.quote_ident(&entity.table_name),
74            def_without_not_null
75        ))
76    }
77}
78
79#[derive(Debug)]
80pub enum MutationExecutorError {
81    Sqlite(rusqlite::Error),
82    SqlCompile(SqlCompileError),
83    UnsupportedValue(&'static str),
84    UnsupportedColumnType(String),
85    Bind(String),
86    Lock(String),
87}
88
89impl std::fmt::Display for MutationExecutorError {
90    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
91        match self {
92            Self::Sqlite(err) => err.fmt(f),
93            Self::SqlCompile(err) => err.fmt(f),
94            Self::UnsupportedValue(kind) => {
95                write!(
96                    f,
97                    "unsupported rusqlite bind value for mutation executor: {kind}"
98                )
99            }
100            Self::UnsupportedColumnType(kind) => {
101                write!(
102                    f,
103                    "unsupported rusqlite column type for record decoding: {kind}"
104                )
105            }
106            Self::Bind(message) => write!(f, "rusqlite bind error: {message}"),
107            Self::Lock(message) => write!(f, "rusqlite connection lock error: {message}"),
108        }
109    }
110}
111
112impl std::error::Error for MutationExecutorError {}
113
114impl From<rusqlite::Error> for MutationExecutorError {
115    fn from(value: rusqlite::Error) -> Self {
116        Self::Sqlite(value)
117    }
118}
119
120impl From<SqlCompileError> for MutationExecutorError {
121    fn from(value: SqlCompileError) -> Self {
122        Self::SqlCompile(value)
123    }
124}
125
126#[derive(Clone)]
127pub struct SqliteMutationExecutor {
128    connection: Arc<Mutex<Connection>>,
129}
130
131impl SqliteMutationExecutor {
132    pub fn new(connection: Arc<Mutex<Connection>>) -> Self {
133        Self { connection }
134    }
135
136    pub fn from_connection(connection: Connection) -> Self {
137        Self::new(Arc::new(Mutex::new(connection)))
138    }
139
140    pub fn connection(&self) -> Arc<Mutex<Connection>> {
141        Arc::clone(&self.connection)
142    }
143
144    pub fn ensure_schema(
145        &self,
146        dialect: &SqliteDialect,
147        entities: &[&EntityDescriptor],
148    ) -> Result<(), MutationExecutorError> {
149        self.ensure_id_space_table(DEFAULT_ID_SPACE_TABLE)?;
150
151        for entity in entities {
152            if !self.table_exists(&entity.table_name)? {
153                let sql = dialect.compile_create_table(entity)?;
154                self.lock()?.execute(&sql, [])?;
155                continue;
156            }
157
158            let existing_columns = self.table_columns(&entity.table_name)?;
159            for property in &entity.properties {
160                let bare_column = strip_identifier_quotes(&property.column_name).to_lowercase();
161                if existing_columns.contains(&bare_column) {
162                    continue;
163                }
164                let sql = dialect.compile_add_column(entity, property)?;
165                self.lock()?.execute(&sql, [])?;
166            }
167
168            for sql in dialect.schema_indexes_sqls(entity)? {
169                self.lock()?.execute(&sql, [])?;
170            }
171        }
172        Ok(())
173    }
174
175    pub fn ensure_id_space_table(&self, table_name: &str) -> Result<(), MutationExecutorError> {
176        let sql = format!(
177            "CREATE TABLE IF NOT EXISTS {} (type_name VARCHAR(100) PRIMARY KEY, current_level BIGINT NOT NULL)",
178            quote_ident(table_name)
179        );
180        self.lock()?.execute(&sql, [])?;
181        Ok(())
182    }
183
184    pub fn begin_transaction(&self) -> Result<(), MutationExecutorError> {
185        self.lock()?.execute("BEGIN IMMEDIATE", [])?;
186        Ok(())
187    }
188
189    pub fn commit_transaction(&self) -> Result<(), MutationExecutorError> {
190        self.lock()?.execute("COMMIT", [])?;
191        Ok(())
192    }
193
194    pub fn rollback_transaction(&self) -> Result<(), MutationExecutorError> {
195        self.lock()?.execute("ROLLBACK", [])?;
196        Ok(())
197    }
198
199    pub fn execute(&self, query: &CompiledQuery) -> Result<u64, MutationExecutorError> {
200        let params = bind_values(&query.params)?;
201        let rows = self
202            .lock()?
203            .execute(&query.sql_with_comment(), params_from_iter(params.iter()))?;
204        Ok(rows as u64)
205    }
206
207    pub fn fetch_all(&self, query: &CompiledQuery) -> Result<Vec<Record>, MutationExecutorError> {
208        let params = bind_values(&query.params)?;
209        let connection = self.lock()?;
210        let mut statement = connection.prepare(&query.sql_with_comment())?;
211        let columns = statement_columns(&statement);
212        let mut rows = statement.query(params_from_iter(params.iter()))?;
213        let mut records = Vec::new();
214        while let Some(row) = rows.next()? {
215            records.push(decode_sqlite_row(row, &columns)?);
216        }
217        Ok(records)
218    }
219
220    /// Fetch rows in streaming mode (chunked).
221    /// Returns a Vec of StreamChunk, each containing up to `chunk_size` rows.
222    pub fn fetch_stream(
223        &self,
224        query: &CompiledQuery,
225        chunk_size: usize,
226    ) -> Result<Vec<teaql_data_service::StreamChunk>, MutationExecutorError> {
227        let params = bind_values(&query.params)?;
228        let connection = self.lock()?;
229        let mut statement = connection.prepare(&query.sql_with_comment())?;
230        let columns = statement_columns(&statement);
231        let mut rows = statement.query(params_from_iter(params.iter()))?;
232
233        let mut chunks = Vec::new();
234        let mut current_chunk = Vec::new();
235        let mut chunk_index = 0;
236
237        while let Some(row) = rows.next()? {
238            current_chunk.push(decode_sqlite_row(row, &columns)?);
239            if current_chunk.len() >= chunk_size {
240                chunks.push(teaql_data_service::StreamChunk {
241                    rows: current_chunk,
242                    chunk_index,
243                    is_last: false,
244                });
245                current_chunk = Vec::new();
246                chunk_index += 1;
247            }
248        }
249
250        // Push the final chunk (may be empty if exactly aligned)
251        chunks.push(teaql_data_service::StreamChunk {
252            rows: current_chunk,
253            chunk_index,
254            is_last: true,
255        });
256
257        Ok(chunks)
258    }
259
260    pub fn table_exists(&self, table_name: &str) -> Result<bool, MutationExecutorError> {
261        let exists: i64 = self.lock()?.query_row(
262            "SELECT COUNT(1) FROM sqlite_master WHERE type = 'table' AND name = ?",
263            [table_name],
264            |row| row.get(0),
265        )?;
266        Ok(exists > 0)
267    }
268
269    pub fn table_columns(
270        &self,
271        table_name: &str,
272    ) -> Result<BTreeSet<String>, MutationExecutorError> {
273        let pragma_sql = format!("PRAGMA table_info({})", quote_ident(table_name));
274        let connection = self.lock()?;
275        let mut statement = connection.prepare(&pragma_sql)?;
276        let rows = statement.query_map([], |row| row.get::<_, String>("name"))?;
277        let mut columns = BTreeSet::new();
278        for row in rows {
279            columns.insert(row?.to_lowercase());
280        }
281        Ok(columns)
282    }
283
284    fn lock(&self) -> Result<MutexGuard<'_, Connection>, MutationExecutorError> {
285        self.connection
286            .lock()
287            .map_err(|err| MutationExecutorError::Lock(err.to_string()))
288    }
289}
290
291impl teaql_data_service::DataServiceExecutor for SqliteMutationExecutor {
292    type Error = MutationExecutorError;
293
294    fn capabilities(&self) -> teaql_data_service::DataServiceCapabilities {
295        teaql_data_service::DataServiceCapabilities {
296            query: true,
297            mutation: true,
298            transaction: true,
299            schema: true,
300            id_generation: true,
301            ..Default::default()
302        }
303    }
304}
305
306impl SqlTransport for SqliteMutationExecutor {
307    type Error = MutationExecutorError;
308
309    async fn fetch_all_sql(&self, query: &CompiledQuery) -> Result<Vec<Record>, Self::Error> {
310        SqliteMutationExecutor::fetch_all(self, query)
311    }
312
313    async fn execute_sql(&self, query: &CompiledQuery) -> Result<u64, Self::Error> {
314        SqliteMutationExecutor::execute(self, query)
315    }
316}
317
318impl teaql_sql::StreamingSqlTransport for SqliteMutationExecutor {
319    fn stream_sql(
320        &self,
321        query: CompiledQuery,
322        chunk_size: usize,
323    ) -> teaql_data_service::QueryStream<'_, Self::Error> {
324        let connection = self.connection.clone();
325        Box::pin(async_stream::try_stream! {
326            let params = bind_values(&query.params)?;
327            let guard = connection.lock().map_err(|err| MutationExecutorError::Lock(err.to_string()))?;
328            let mut statement = guard.prepare(&query.sql_with_comment())?;
329            let columns = statement_columns(&statement);
330            let mut rows = statement.query(params_from_iter(params.iter()))?;
331            let mut chunk = Vec::with_capacity(chunk_size); let mut index = 0;
332            while let Some(row) = rows.next()? {
333                chunk.push(decode_sqlite_row(row, &columns)?);
334                if chunk.len() == chunk_size { yield teaql_data_service::StreamChunk { rows: std::mem::take(&mut chunk), chunk_index: index, is_last: false }; index += 1; }
335            }
336            if !chunk.is_empty() { yield teaql_data_service::StreamChunk { rows: chunk, chunk_index: index, is_last: true }; }
337        })
338    }
339}
340
341impl teaql_data_service::StreamQueryExecutor for SqliteMutationExecutor {
342    fn query_stream(
343        &self,
344        request: teaql_data_service::QueryRequest,
345        chunk_size: usize,
346    ) -> teaql_data_service::QueryStream<'_, Self::Error> {
347        let dialect = SqliteDialect;
348        // Use a dummy entity descriptor for compilation
349        let entity_desc = teaql_core::EntityDescriptor::new(&request.query.entity);
350        match dialect.compile_select(&entity_desc, &request.query) {
351            Ok(compiled) => {
352                teaql_sql::StreamingSqlTransport::stream_sql(self, compiled, chunk_size)
353            }
354            Err(error) => Box::pin(futures_util::stream::once(async {
355                Err(MutationExecutorError::SqlCompile(error))
356            })),
357        }
358    }
359}
360
361impl teaql_sql::SqlTransaction for SqliteMutationExecutor {
362    type Error = MutationExecutorError;
363
364    async fn commit_sql(self) -> Result<(), Self::Error> {
365        self.commit_transaction()
366    }
367
368    async fn rollback_sql(self) -> Result<(), Self::Error> {
369        self.rollback_transaction()
370    }
371}
372
373impl teaql_sql::SqlTransactionTransport for SqliteMutationExecutor {
374    type Tx<'a>
375        = Self
376    where
377        Self: 'a;
378
379    async fn begin_sql(&self) -> Result<Self::Tx<'_>, Self::Error> {
380        self.begin_transaction()?;
381        Ok(self.clone())
382    }
383}
384
385fn initial_graph_exists_sqlite(
386    executor: &SqliteMutationExecutor,
387    dialect: &SqliteDialect,
388    entity: &EntityDescriptor,
389    graph: &GraphNode,
390) -> Result<bool, MutationExecutorError> {
391    let Some(id) = graph.values.get("id") else {
392        return Ok(false);
393    };
394    let query = dialect.compile_select(
395        entity,
396        &SelectQuery::new(&graph.entity)
397            .project("id")
398            .filter(Expr::eq("id", id.clone()))
399            .limit(1),
400    )?;
401    Ok(!executor.fetch_all(&query)?.is_empty())
402}
403
404fn compile_initial_graph_insert(
405    dialect: &impl SqlDialect,
406    entity: &EntityDescriptor,
407    graph: &GraphNode,
408) -> Result<CompiledQuery, MutationExecutorError> {
409    let mut command = InsertCommand::new(&graph.entity);
410    for (field, value) in &graph.values {
411        command = command.value(field.clone(), value.clone());
412    }
413    dialect.compile_insert(entity, &command).map_err(Into::into)
414}
415
416fn compile_initial_graph_update(
417    dialect: &impl SqlDialect,
418    entity: &EntityDescriptor,
419    graph: &GraphNode,
420) -> Result<Option<CompiledQuery>, MutationExecutorError> {
421    let Some(id) = graph.values.get("id") else {
422        return Ok(None);
423    };
424    let mut command = UpdateCommand::new(&graph.entity, id.clone());
425    for (field, value) in &graph.values {
426        if field == "id" {
427            continue;
428        }
429        command = command.value(field.clone(), value.clone());
430    }
431    match dialect.compile_update(entity, &command) {
432        Ok(query) => Ok(Some(query)),
433        Err(SqlCompileError::EmptyMutation(_)) => Ok(None),
434        Err(err) => Err(err.into()),
435    }
436}
437
438pub trait SqliteSchemaExt {
439    fn ensure_sqlite_schema(
440        &self,
441    ) -> Pin<Box<dyn Future<Output = Result<(), MutationExecutorError>> + Send + '_>>;
442}
443
444pub fn ensure_sqlite_schema_for(ctx: &UserContext) -> Result<(), MutationExecutorError> {
445    let dialect = ctx.get_resource::<SqliteDialect>().ok_or_else(|| {
446        MutationExecutorError::Bind("missing typed resource: SqliteDialect".to_owned())
447    })?;
448    let executor = ctx
449        .get_resource::<SqliteMutationExecutor>()
450        .ok_or_else(|| {
451            MutationExecutorError::Bind("missing typed resource: SqliteMutationExecutor".to_owned())
452        })?;
453
454    let entities = ctx.all_entities();
455
456    // Ensure id space table exists
457    executor.ensure_id_space_table(DEFAULT_ID_SPACE_TABLE)?;
458
459    // Process each entity table individually with granular events
460    for entity in &entities {
461        let field_count = entity.properties.len();
462        if !executor.table_exists(&entity.table_name)? {
463            // New table: create it
464            let sql = dialect.compile_create_table(entity)?;
465            executor.lock()?.execute(&sql, [])?;
466            let _ = ctx.send_event(RawAuditEvent::schema_created(
467                &entity.name,
468                &entity.table_name,
469                field_count,
470            ));
471            continue;
472        }
473        // Existing table: check for missing columns
474        let existing_columns = executor.table_columns(&entity.table_name)?;
475        let mut fields_added = 0;
476        for property in &entity.properties {
477            let bare_column = strip_identifier_quotes(&property.column_name).to_lowercase();
478            if existing_columns.contains(&bare_column) {
479                continue;
480            }
481            let sql = dialect.compile_add_column(entity, property)?;
482            executor.lock()?.execute(&sql, [])?;
483            let _ = ctx.send_event(RawAuditEvent::field_added(
484                &entity.name,
485                &entity.table_name,
486                &property.column_name,
487            ));
488            fields_added += 1;
489        }
490        let _ = ctx.send_event(RawAuditEvent::schema_verified(
491            &entity.name,
492            &entity.table_name,
493            field_count,
494        ));
495        let _ = fields_added; // used above for FieldAdded events
496    }
497
498    // Seed initial data, tracking insert vs update counts per entity
499    let mut seed_counts: BTreeMap<String, (usize, usize)> = BTreeMap::new(); // (inserted, updated)
500    for graph in ctx.initial_graphs() {
501        let entity = ctx.entity(&graph.entity).ok_or_else(|| {
502            MutationExecutorError::Bind(format!("missing entity: {}", graph.entity))
503        })?;
504        let counts = seed_counts.entry(graph.entity.clone()).or_insert((0, 0));
505        if initial_graph_exists_sqlite(executor, dialect, entity, graph)? {
506            if let Some(query) = compile_initial_graph_update(dialect, entity, graph)? {
507                executor.execute(&query)?;
508            }
509            counts.1 += 1; // updated
510            continue;
511        }
512        let query = compile_initial_graph_insert(dialect, entity, graph)?;
513        executor.execute(&query)?;
514        counts.0 += 1; // inserted
515    }
516
517    // Fire DataSeeded events per entity type
518    for (entity_name, (inserted, updated)) in &seed_counts {
519        let entity = ctx.entity(entity_name).ok_or_else(|| {
520            MutationExecutorError::Bind(format!("missing entity: {}", entity_name))
521        })?;
522        let _ = ctx.send_event(RawAuditEvent::data_seeded(
523            entity_name,
524            &entity.table_name,
525            *inserted,
526            *updated,
527        ));
528    }
529
530    Ok(())
531}
532
533impl SqliteSchemaExt for UserContext {
534    fn ensure_sqlite_schema(
535        &self,
536    ) -> Pin<Box<dyn Future<Output = Result<(), MutationExecutorError>> + Send + '_>> {
537        Box::pin(async move { ensure_sqlite_schema_for(self) })
538    }
539}
540
541#[derive(Debug, Default, Clone, Copy)]
542pub struct SqliteSchemaProvider;
543
544impl SchemaProvider for SqliteSchemaProvider {
545    fn ensure_schema<'a>(
546        &'a self,
547        ctx: &'a UserContext,
548    ) -> Pin<Box<dyn Future<Output = Result<(), RuntimeError>> + Send + 'a>> {
549        Box::pin(async move {
550            ensure_sqlite_schema_for(ctx).map_err(|err| RuntimeError::Schema(err.to_string()))
551        })
552    }
553}
554
555pub trait SqliteProviderExt {
556    fn use_sqlite_provider(&mut self, executor: SqliteMutationExecutor) -> &mut Self;
557}
558
559impl SqliteProviderExt for UserContext {
560    fn use_sqlite_provider(&mut self, executor: SqliteMutationExecutor) -> &mut Self {
561        self.insert_resource(SqliteDialect);
562        self.insert_resource(executor);
563        self.set_schema_provider(SqliteSchemaProvider);
564        self
565    }
566}
567
568#[derive(Clone)]
569pub struct SqliteIdSpaceGenerator {
570    executor: SqliteMutationExecutor,
571    table_name: String,
572}
573
574impl SqliteIdSpaceGenerator {
575    pub fn new(connection: Connection) -> Self {
576        Self::from_executor(SqliteMutationExecutor::from_connection(connection))
577    }
578
579    pub fn from_executor(executor: SqliteMutationExecutor) -> Self {
580        Self {
581            executor,
582            table_name: DEFAULT_ID_SPACE_TABLE.to_owned(),
583        }
584    }
585
586    pub fn with_table_name(mut self, table_name: impl Into<String>) -> Self {
587        self.table_name = table_name.into();
588        self
589    }
590
591    pub fn ensure_table(&self) -> Result<(), MutationExecutorError> {
592        self.executor.ensure_id_space_table(&self.table_name)
593    }
594
595    pub fn next_id(&self, entity: &str) -> Result<u64, MutationExecutorError> {
596        self.ensure_table()?;
597        let sql = format!(
598            "INSERT INTO {} (type_name, current_level) VALUES (?, 1) \
599             ON CONFLICT (type_name) DO UPDATE \
600             SET current_level = current_level + 1 \
601             RETURNING current_level",
602            quote_ident(&self.table_name)
603        );
604        let id: i64 = self
605            .executor
606            .lock()?
607            .query_row(&sql, [entity], |row| row.get(0))?;
608        u64::try_from(id).map_err(|_| {
609            MutationExecutorError::Bind(format!("generated id {id} cannot be represented as u64"))
610        })
611    }
612}
613
614impl InternalIdGenerator for SqliteIdSpaceGenerator {
615    fn generate_id(&self, entity: &str) -> Result<u64, RuntimeError> {
616        self.next_id(entity)
617            .map_err(|err| RuntimeError::IdGeneration(err.to_string()))
618    }
619}
620
621fn quote_ident(ident: &str) -> String {
622    quote_identifier_if_needed(ident, '"')
623}
624
625/// Strip wrapping identifier quotes from a SQL identifier.
626///
627/// SQLite `PRAGMA table_info` returns bare column names (e.g. `description`),
628/// but generated `PropertyDescriptor::column_name` may carry quotes
629/// (e.g. `"description"`) when the name is a reserved keyword.  This helper
630/// normalises the column name so the two can be compared correctly during
631/// schema migration.
632fn strip_identifier_quotes(ident: &str) -> &str {
633    let bytes = ident.as_bytes();
634    if bytes.len() >= 2 {
635        let (first, last) = (bytes[0], bytes[bytes.len() - 1]);
636        if (first == b'"' && last == b'"')
637            || (first == b'`' && last == b'`')
638            || (first == b'[' && last == b']')
639        {
640            return &ident[1..ident.len() - 1];
641        }
642    }
643    ident
644}
645
646fn bind_values(values: &[Value]) -> Result<Vec<SqliteValue>, MutationExecutorError> {
647    values.iter().map(bind_sqlite_value).collect()
648}
649
650fn bind_sqlite_value(value: &Value) -> Result<SqliteValue, MutationExecutorError> {
651    match value {
652        Value::Null => Ok(SqliteValue::Null),
653        Value::Bool(v) => Ok(SqliteValue::Integer(i64::from(*v))),
654        Value::I64(v) => Ok(SqliteValue::Integer(*v)),
655        Value::U64(v) => i64::try_from(*v)
656            .map(SqliteValue::Integer)
657            .map_err(|_| MutationExecutorError::Bind(format!("u64 value {v} exceeds i64 range"))),
658        Value::F64(v) => Ok(SqliteValue::Real(*v)),
659        // Bind the canonical numeric spelling. SQLite NUMERIC affinity keeps
660        // predicates and aggregates numeric; an application-only text prefix
661        // makes range comparisons silently return the wrong result.
662        Value::Decimal(v) => Ok(SqliteValue::Text(v.to_string())),
663        Value::Text(v) => Ok(SqliteValue::Text(v.clone())),
664        Value::Json(v) => Ok(SqliteValue::Text(v.to_string())),
665        Value::Date(v) => Ok(SqliteValue::Text(v.format("%Y-%m-%d").to_string())),
666        Value::Timestamp(v) => Ok(SqliteValue::Text(v.0.to_string())),
667        Value::Object(_) => Err(MutationExecutorError::UnsupportedValue("object")),
668        Value::List(_) => Err(MutationExecutorError::UnsupportedValue("list")),
669        Value::TypedNull(_) => Ok(SqliteValue::Null),
670    }
671}
672
673#[derive(Debug, Clone)]
674struct ColumnInfo {
675    name: String,
676    decl_type: Option<String>,
677}
678
679fn statement_columns(statement: &rusqlite::Statement<'_>) -> Vec<ColumnInfo> {
680    statement
681        .columns()
682        .into_iter()
683        .map(|column| ColumnInfo {
684            name: column.name().to_owned(),
685            decl_type: column.decl_type().map(|value| value.to_ascii_uppercase()),
686        })
687        .collect()
688}
689
690fn decode_sqlite_row(
691    row: &Row<'_>,
692    columns: &[ColumnInfo],
693) -> Result<Record, MutationExecutorError> {
694    let mut record = BTreeMap::new();
695    for (index, column) in columns.iter().enumerate() {
696        let value_ref = row.get_ref(index)?;
697        let value = match value_ref {
698            ValueRef::Null => Value::Null,
699            ValueRef::Integer(value) => decode_sqlite_integer(value, column),
700            ValueRef::Real(value) => Value::F64(value),
701            ValueRef::Text(value) => decode_sqlite_text(value, column)?,
702            ValueRef::Blob(_) => {
703                return Err(MutationExecutorError::UnsupportedColumnType(
704                    "BLOB".to_owned(),
705                ));
706            }
707        };
708        record.insert(column.name.clone(), value);
709    }
710    Ok(record)
711}
712
713fn decode_sqlite_integer(value: i64, column: &ColumnInfo) -> Value {
714    match column_decl_type(column).as_deref() {
715        Some("BOOLEAN") | Some("BOOL") => Value::Bool(value != 0),
716        _ => Value::I64(value),
717    }
718}
719
720fn decode_sqlite_text(value: &[u8], column: &ColumnInfo) -> Result<Value, MutationExecutorError> {
721    let value = std::str::from_utf8(value)
722        .map_err(|err| MutationExecutorError::Bind(format!("invalid sqlite text: {err}")))?;
723    match column_decl_type(column).as_deref() {
724        Some("NUMERIC") | Some("DECIMAL") => Decimal::from_str(value)
725            .map(Value::Decimal)
726            .map_err(|err| MutationExecutorError::Bind(format!("invalid sqlite decimal: {err}"))),
727        Some("JSON") => serde_json::from_str(value).map(Value::Json).map_err(|err| {
728            MutationExecutorError::Bind(format!("invalid sqlite json value: {err}"))
729        }),
730        Some("DATE") => NaiveDate::parse_from_str(value, "%Y-%m-%d")
731            .map(Value::Date)
732            .map_err(|err| MutationExecutorError::Bind(format!("invalid sqlite date: {err}"))),
733        Some("TIMESTAMP") | Some("DATETIME") => parse_sqlite_timestamp(value),
734        _ => infer_sqlite_text(value),
735    }
736}
737
738fn infer_sqlite_text(value: &str) -> Result<Value, MutationExecutorError> {
739    if let Ok(date) = NaiveDate::parse_from_str(value, "%Y-%m-%d") {
740        return Ok(Value::Date(date));
741    }
742    if let Ok(timestamp) = DateTime::parse_from_rfc3339(value) {
743        return Ok(Value::Timestamp(teaql_core::time::Timestamp(
744            timestamp.timestamp_millis(),
745        )));
746    }
747    if let Ok(timestamp) = NaiveDateTime::parse_from_str(value, "%Y-%m-%d %H:%M:%S") {
748        return Ok(Value::Timestamp(teaql_core::time::Timestamp(
749            timestamp.and_utc().timestamp_millis(),
750        )));
751    }
752    Ok(Value::Text(value.to_owned()))
753}
754
755fn parse_sqlite_timestamp(value: &str) -> Result<Value, MutationExecutorError> {
756    if let Ok(timestamp) = DateTime::parse_from_rfc3339(value) {
757        return Ok(Value::Timestamp(teaql_core::time::Timestamp(
758            timestamp.timestamp_millis(),
759        )));
760    }
761    if let Ok(date) = NaiveDate::parse_from_str(value, "%Y-%m-%d") {
762        return Ok(Value::Timestamp(teaql_core::time::Timestamp(
763            date.and_hms_opt(0, 0, 0)
764                .unwrap_or_default()
765                .and_utc()
766                .timestamp_millis(),
767        )));
768    }
769    NaiveDateTime::parse_from_str(value, "%Y-%m-%d %H:%M:%S")
770        .map(|timestamp| {
771            Value::Timestamp(teaql_core::time::Timestamp(
772                timestamp.and_utc().timestamp_millis(),
773            ))
774        })
775        .map_err(|err| MutationExecutorError::Bind(format!("invalid sqlite timestamp: {err}")))
776}
777
778fn column_decl_type(column: &ColumnInfo) -> Option<String> {
779    column
780        .decl_type
781        .as_ref()
782        .map(|value| value.split('(').next().unwrap_or(value).trim().to_owned())
783}
784
785#[cfg(test)]
786mod tests {
787    use super::*;
788    use futures_util::StreamExt;
789    use teaql_core::{DeleteCommand, RecoverCommand};
790    use teaql_macros::TeaqlEntity;
791    use teaql_runtime::InMemoryMetadataStore;
792
793    #[test]
794    fn streaming_sql_yields_bounded_chunks_and_releases_cursor_on_drop() {
795        let connection = Connection::open_in_memory().unwrap();
796        connection
797            .execute_batch(
798                "CREATE TABLE stream_fixture(id INTEGER);\
799                 INSERT INTO stream_fixture VALUES (1), (2), (3), (4), (5);",
800            )
801            .unwrap();
802        let executor = SqliteMutationExecutor::from_connection(connection);
803        let query = CompiledQuery {
804            sql: "SELECT id FROM stream_fixture ORDER BY id".to_owned(),
805            params: vec![],
806            comment: None,
807        };
808        let mut stream = teaql_sql::StreamingSqlTransport::stream_sql(&executor, query.clone(), 2);
809        let sizes = futures_executor::block_on(async {
810            let mut result = Vec::new();
811            while let Some(chunk) = stream.next().await {
812                result.push(chunk.unwrap().rows.len());
813            }
814            result
815        });
816        assert_eq!(sizes, vec![2, 2, 1]);
817
818        let mut early = teaql_sql::StreamingSqlTransport::stream_sql(&executor, query, 2);
819        assert_eq!(
820            futures_executor::block_on(early.next())
821                .unwrap()
822                .unwrap()
823                .rows
824                .len(),
825            2
826        );
827        drop(early);
828        let count: i64 = executor
829            .connection()
830            .lock()
831            .unwrap()
832            .query_row("SELECT count(*) FROM stream_fixture", [], |row| row.get(0))
833            .unwrap();
834        assert_eq!(count, 5);
835    }
836
837    #[test]
838    fn decimal_bind_is_numeric_and_comparable() {
839        let value =
840            bind_sqlite_value(&Value::Decimal(Decimal::from_str("123.450").unwrap())).unwrap();
841        assert_eq!(value, SqliteValue::Text("123.450".to_owned()));
842        let connection = Connection::open_in_memory().unwrap();
843        let matches: i64 = connection
844            .query_row(
845                "SELECT 1 WHERE CAST(? AS NUMERIC) BETWEEN 120 AND 130",
846                [value],
847                |row| row.get(0),
848            )
849            .unwrap();
850        assert_eq!(matches, 1);
851    }
852
853    fn entity() -> EntityDescriptor {
854        EntityDescriptor::new("Order")
855            .table_name("orders")
856            .property(
857                PropertyDescriptor::new("id", DataType::U64)
858                    .column_name("id")
859                    .id()
860                    .not_null(),
861            )
862            .property(
863                PropertyDescriptor::new("version", DataType::I64)
864                    .column_name("version")
865                    .version()
866                    .not_null(),
867            )
868            .property(PropertyDescriptor::new("name", DataType::Text).column_name("name"))
869    }
870
871    fn order_line_entity() -> EntityDescriptor {
872        EntityDescriptor::new("OrderLine")
873            .table_name("order_line")
874            .property(
875                PropertyDescriptor::new("id", DataType::U64)
876                    .column_name("id")
877                    .id()
878                    .not_null(),
879            )
880            .property(
881                PropertyDescriptor::new("order_id", DataType::U64)
882                    .column_name("order_id")
883                    .not_null(),
884            )
885            .property(PropertyDescriptor::new("name", DataType::Text).column_name("name"))
886    }
887
888    #[allow(dead_code)]
889    #[derive(Debug, PartialEq, TeaqlEntity)]
890    #[teaql(entity = "FeatureFlag", table = "feature_flags")]
891    struct FeatureFlagRow {
892        #[teaql(id)]
893        id: u64,
894        #[teaql(version)]
895        version: i64,
896        enabled: bool,
897        optional_enabled: Option<bool>,
898    }
899
900    fn feature_flag_record(enabled: Value, optional_enabled: Value) -> Record {
901        Record::from([
902            ("id".to_owned(), Value::U64(1)),
903            ("version".to_owned(), Value::I64(1)),
904            ("enabled".to_owned(), enabled),
905            ("optional_enabled".to_owned(), optional_enabled),
906        ])
907    }
908
909    #[test]
910    fn sqlite_dialect_compiles_mutations_and_schema() {
911        let insert = SqliteDialect
912            .compile_insert(
913                &entity(),
914                &InsertCommand::new("Order")
915                    .value("id", 1_u64)
916                    .value("name", "A"),
917            )
918            .unwrap();
919        assert_eq!(insert.sql, "INSERT INTO orders (id, name) VALUES (?, ?)");
920
921        let update = SqliteDialect
922            .compile_update(
923                &entity(),
924                &UpdateCommand::new("Order", 1_u64)
925                    .expected_version(3)
926                    .value("name", "B"),
927            )
928            .unwrap();
929        assert_eq!(
930            update.sql,
931            "UPDATE orders SET name = ?, version = ? WHERE id = ? AND version = ?"
932        );
933
934        let delete = SqliteDialect
935            .compile_delete(
936                &entity(),
937                &DeleteCommand::new("Order", 1_u64).expected_version(3),
938            )
939            .unwrap();
940        let recover = SqliteDialect
941            .compile_recover(&entity(), &RecoverCommand::new("Order", 1_u64, -4))
942            .unwrap();
943        assert_eq!(
944            delete.sql,
945            "UPDATE orders SET version = ? WHERE id = ? AND version = ?"
946        );
947        assert_eq!(
948            recover.sql,
949            "UPDATE orders SET version = ? WHERE id = ? AND version = ?"
950        );
951
952        let create = SqliteDialect.compile_create_table(&entity()).unwrap();
953        assert_eq!(
954            create,
955            "CREATE TABLE IF NOT EXISTS orders (id INTEGER PRIMARY KEY NOT NULL, version INTEGER NOT NULL, name VARCHAR(255))"
956        );
957    }
958
959    #[test]
960    fn sqlite_executor_ensures_schema_and_roundtrips_rows() {
961        let executor =
962            SqliteMutationExecutor::from_connection(Connection::open_in_memory().unwrap());
963        let entity = entity();
964        let mut ctx = UserContext::new()
965            .with_metadata(InMemoryMetadataStore::new().with_entity(entity.clone()));
966
967        ctx.use_sqlite_provider(executor.clone());
968        ensure_sqlite_schema_for(&ctx).unwrap();
969
970        let insert = SqliteDialect
971            .compile_insert(
972                &entity,
973                &InsertCommand::new("Order")
974                    .value("id", 1_u64)
975                    .value("version", 1_i64)
976                    .value("name", "draft"),
977            )
978            .unwrap();
979        assert_eq!(executor.execute(&insert).unwrap(), 1);
980
981        let select = SqliteDialect
982            .compile_select(
983                &entity,
984                &SelectQuery::new("Order")
985                    .filter(Expr::eq("id", 1_u64))
986                    .order_asc("id"),
987            )
988            .unwrap();
989        let rows = executor.fetch_all(&select).unwrap();
990        assert_eq!(rows.len(), 1);
991        assert_eq!(rows[0].get("id"), Some(&Value::I64(1)));
992        assert_eq!(rows[0].get("version"), Some(&Value::I64(1)));
993        assert_eq!(rows[0].get("name"), Some(&Value::Text("draft".to_owned())));
994    }
995
996    #[test]
997    fn sqlite_executes_partitioned_relation_limit_per_parent() {
998        let executor =
999            SqliteMutationExecutor::from_connection(Connection::open_in_memory().unwrap());
1000        let entity = order_line_entity();
1001        executor.ensure_schema(&SqliteDialect, &[&entity]).unwrap();
1002
1003        for order_id in [11_u64, 12_u64] {
1004            for index in 1_u64..=5 {
1005                let id = order_id * 100 + index;
1006                let insert = SqliteDialect
1007                    .compile_insert(
1008                        &entity,
1009                        &InsertCommand::new("OrderLine")
1010                            .value("id", id)
1011                            .value("order_id", order_id)
1012                            .value("name", format!("line-{id}")),
1013                    )
1014                    .unwrap();
1015                executor.execute(&insert).unwrap();
1016            }
1017        }
1018
1019        let query = SelectQuery::new("OrderLine")
1020            .project("id")
1021            .project("order_id")
1022            .order_desc("id")
1023            .limit(3)
1024            .partition_by("order_id");
1025        let compiled = SqliteDialect.compile_select(&entity, &query).unwrap();
1026        let rows = executor.fetch_all(&compiled).unwrap();
1027
1028        assert_eq!(rows.len(), 6);
1029        for order_id in [11_i64, 12_i64] {
1030            let ids = rows
1031                .iter()
1032                .filter(|row| row.get("order_id") == Some(&Value::I64(order_id)))
1033                .filter_map(|row| row.get("id").cloned())
1034                .collect::<Vec<_>>();
1035            assert_eq!(
1036                ids,
1037                vec![
1038                    Value::I64(order_id * 100 + 5),
1039                    Value::I64(order_id * 100 + 4),
1040                    Value::I64(order_id * 100 + 3),
1041                ]
1042            );
1043        }
1044    }
1045
1046    #[test]
1047    fn sqlite_boolean_new_schema_roundtrips_as_bool() {
1048        let executor =
1049            SqliteMutationExecutor::from_connection(Connection::open_in_memory().unwrap());
1050        let entity = <FeatureFlagRow as teaql_core::TeaqlEntity>::entity_descriptor();
1051        let ddl = SqliteDialect.compile_create_table(&entity).unwrap();
1052        assert!(ddl.contains("enabled BOOLEAN NOT NULL"), "{ddl}");
1053        assert!(ddl.contains("optional_enabled BOOLEAN"), "{ddl}");
1054        assert!(!ddl.contains("enabled INTEGER"), "{ddl}");
1055
1056        executor.ensure_schema(&SqliteDialect, &[&entity]).unwrap();
1057        for (id, enabled, optional_enabled) in [(1_u64, false, true), (2_u64, true, false)] {
1058            let insert = SqliteDialect
1059                .compile_insert(
1060                    &entity,
1061                    &InsertCommand::new("FeatureFlag")
1062                        .value("id", id)
1063                        .value("version", 1_i64)
1064                        .value("enabled", enabled)
1065                        .value("optional_enabled", optional_enabled),
1066                )
1067                .unwrap();
1068            assert_eq!(executor.execute(&insert).unwrap(), 1);
1069        }
1070
1071        let select = SqliteDialect
1072            .compile_select(&entity, &SelectQuery::new("FeatureFlag").order_asc("id"))
1073            .unwrap();
1074        let rows = executor.fetch_all(&select).unwrap();
1075        assert_eq!(rows[0].get("enabled"), Some(&Value::Bool(false)));
1076        assert_eq!(rows[0].get("optional_enabled"), Some(&Value::Bool(true)));
1077        assert_eq!(rows[1].get("enabled"), Some(&Value::Bool(true)));
1078        assert_eq!(rows[1].get("optional_enabled"), Some(&Value::Bool(false)));
1079
1080        let first = <FeatureFlagRow as teaql_core::Entity>::from_record(rows[0].clone()).unwrap();
1081        let second = <FeatureFlagRow as teaql_core::Entity>::from_record(rows[1].clone()).unwrap();
1082        assert!(!first.enabled);
1083        assert_eq!(first.optional_enabled, Some(true));
1084        assert!(second.enabled);
1085        assert_eq!(second.optional_enabled, Some(false));
1086    }
1087
1088    #[test]
1089    fn sqlite_boolean_legacy_integer_schema_maps_only_binary_values() {
1090        let executor =
1091            SqliteMutationExecutor::from_connection(Connection::open_in_memory().unwrap());
1092        let entity = <FeatureFlagRow as teaql_core::TeaqlEntity>::entity_descriptor();
1093        executor
1094            .execute(&CompiledQuery {
1095                sql: "CREATE TABLE feature_flags (id INTEGER PRIMARY KEY, version INTEGER NOT NULL, enabled INTEGER NOT NULL, optional_enabled INTEGER)"
1096                    .to_owned(),
1097                params: Vec::new(),
1098                comment: None,
1099            })
1100            .unwrap();
1101
1102        let insert = SqliteDialect
1103            .compile_insert(
1104                &entity,
1105                &InsertCommand::new("FeatureFlag")
1106                    .value("id", 1_u64)
1107                    .value("version", 1_i64)
1108                    .value("enabled", true)
1109                    .value("optional_enabled", false),
1110            )
1111            .unwrap();
1112        executor.execute(&insert).unwrap();
1113        executor
1114            .execute(&CompiledQuery {
1115                sql: "INSERT INTO feature_flags (id, version, enabled, optional_enabled) VALUES (?, ?, ?, ?)"
1116                    .to_owned(),
1117                params: vec![
1118                    Value::U64(2),
1119                    Value::I64(1),
1120                    Value::I64(2),
1121                    Value::Null,
1122                ],
1123                comment: None,
1124            })
1125            .unwrap();
1126        let select = SqliteDialect
1127            .compile_select(&entity, &SelectQuery::new("FeatureFlag").order_asc("id"))
1128            .unwrap();
1129        let rows = executor.fetch_all(&select).unwrap();
1130        assert_eq!(rows[0].get("version"), Some(&Value::I64(1)));
1131        assert_eq!(rows[0].get("enabled"), Some(&Value::I64(1)));
1132        assert_eq!(rows[0].get("optional_enabled"), Some(&Value::I64(0)));
1133
1134        let decoded = <FeatureFlagRow as teaql_core::Entity>::from_record(rows[0].clone()).unwrap();
1135        assert!(decoded.enabled);
1136        assert_eq!(decoded.optional_enabled, Some(false));
1137        assert_eq!(rows[1].get("enabled"), Some(&Value::I64(2)));
1138        let error =
1139            <FeatureFlagRow as teaql_core::Entity>::from_record(rows[1].clone()).unwrap_err();
1140        assert!(error.message.contains("invalid field enabled"));
1141
1142        for (value, expected) in [
1143            (Value::I64(0), false),
1144            (Value::I64(1), true),
1145            (Value::U64(0), false),
1146            (Value::U64(1), true),
1147        ] {
1148            let decoded = <FeatureFlagRow as teaql_core::Entity>::from_record(feature_flag_record(
1149                value,
1150                Value::Null,
1151            ))
1152            .unwrap();
1153            assert_eq!(decoded.enabled, expected);
1154            assert_eq!(decoded.optional_enabled, None);
1155        }
1156
1157        for invalid in [Value::I64(-1), Value::I64(2), Value::U64(2)] {
1158            let error = <FeatureFlagRow as teaql_core::Entity>::from_record(feature_flag_record(
1159                invalid,
1160                Value::Null,
1161            ))
1162            .unwrap_err();
1163            assert!(error.message.contains("invalid field enabled"));
1164        }
1165        let error = <FeatureFlagRow as teaql_core::Entity>::from_record(feature_flag_record(
1166            Value::Bool(true),
1167            Value::U64(2),
1168        ))
1169        .unwrap_err();
1170        assert!(error.message.contains("invalid field optional_enabled"));
1171    }
1172
1173    #[test]
1174    fn sqlite_executor_parses_json_only_for_json_columns() {
1175        let executor =
1176            SqliteMutationExecutor::from_connection(Connection::open_in_memory().unwrap());
1177
1178        executor
1179            .execute(&CompiledQuery {
1180                sql: "CREATE TABLE payloads (text_payload TEXT, json_payload JSON)".to_owned(),
1181                params: Vec::new(),
1182                comment: None,
1183            })
1184            .unwrap();
1185        executor
1186            .execute(&CompiledQuery {
1187                sql: "INSERT INTO payloads (text_payload, json_payload) VALUES (?, ?)".to_owned(),
1188                params: vec![
1189                    Value::Text("{\"active\":true}".to_owned()),
1190                    Value::Json(serde_json::json!({"active": true})),
1191                ],
1192                comment: None,
1193            })
1194            .unwrap();
1195
1196        let rows = executor
1197            .fetch_all(&CompiledQuery {
1198                sql: "SELECT text_payload, json_payload FROM payloads".to_owned(),
1199                params: Vec::new(),
1200                comment: None,
1201            })
1202            .unwrap();
1203
1204        assert_eq!(
1205            rows[0].get("text_payload"),
1206            Some(&Value::Text("{\"active\":true}".to_owned()))
1207        );
1208        assert_eq!(
1209            rows[0].get("json_payload"),
1210            Some(&Value::Json(serde_json::json!({"active": true})))
1211        );
1212    }
1213
1214    #[test]
1215    fn sqlite_id_space_generator_increments_ids() {
1216        let executor =
1217            SqliteMutationExecutor::from_connection(Connection::open_in_memory().unwrap());
1218        let generator = SqliteIdSpaceGenerator::from_executor(executor);
1219        assert_eq!(generator.next_id("Order").unwrap(), 1);
1220        assert_eq!(generator.next_id("Order").unwrap(), 2);
1221    }
1222
1223    #[test]
1224    fn sqlite_fetch_stream_returns_chunked_rows() {
1225        let executor = SqliteMutationExecutor::new(Arc::new(Mutex::new(
1226            Connection::open_in_memory().unwrap(),
1227        )));
1228        let entity = entity();
1229
1230        // Create table and insert 25 rows
1231        executor
1232            .execute(&CompiledQuery {
1233                sql: "CREATE TABLE orders (id INTEGER PRIMARY KEY, version INTEGER, name VARCHAR(255))"
1234                    .to_owned(),
1235                params: Vec::new(),
1236                comment: None,
1237            })
1238            .unwrap();
1239
1240        for i in 1..=25 {
1241            let insert = SqliteDialect
1242                .compile_insert(
1243                    &entity,
1244                    &InsertCommand::new("Order")
1245                        .value("id", i as u64)
1246                        .value("version", 1_i64)
1247                        .value("name", format!("order-{i}")),
1248                )
1249                .unwrap();
1250            executor.execute(&insert).unwrap();
1251        }
1252
1253        // Stream with chunk_size = 10
1254        let query = SelectQuery::new("Order")
1255            .filter(Expr::gt("version", 0_i64))
1256            .order_asc("id")
1257            .stream(10);
1258
1259        let compiled = SqliteDialect.compile_select(&entity, &query).unwrap();
1260
1261        let chunks = executor.fetch_stream(&compiled, 10).unwrap();
1262
1263        // 25 rows / 10 per chunk = 3 chunks
1264        assert_eq!(chunks.len(), 3);
1265        assert_eq!(chunks[0].rows.len(), 10);
1266        assert_eq!(chunks[0].chunk_index, 0);
1267        assert!(!chunks[0].is_last);
1268
1269        assert_eq!(chunks[1].rows.len(), 10);
1270        assert_eq!(chunks[1].chunk_index, 1);
1271        assert!(!chunks[1].is_last);
1272
1273        assert_eq!(chunks[2].rows.len(), 5);
1274        assert_eq!(chunks[2].chunk_index, 2);
1275        assert!(chunks[2].is_last);
1276
1277        // Verify first and last row
1278        assert_eq!(
1279            chunks[0].rows[0].get("name"),
1280            Some(&Value::Text("order-1".to_owned()))
1281        );
1282        assert_eq!(
1283            chunks[2].rows[4].get("name"),
1284            Some(&Value::Text("order-25".to_owned()))
1285        );
1286    }
1287
1288    #[test]
1289    fn sqlite_fetch_stream_handles_empty_result() {
1290        let executor = SqliteMutationExecutor::new(Arc::new(Mutex::new(
1291            Connection::open_in_memory().unwrap(),
1292        )));
1293
1294        executor
1295            .execute(&CompiledQuery {
1296                sql: "CREATE TABLE orders (id INTEGER PRIMARY KEY, version INTEGER, name VARCHAR(255))"
1297                    .to_owned(),
1298                params: Vec::new(),
1299                comment: None,
1300            })
1301            .unwrap();
1302
1303        let entity = entity();
1304        let query = SelectQuery::new("Order")
1305            .filter(Expr::gt("version", 0_i64))
1306            .stream(10);
1307
1308        let compiled = SqliteDialect.compile_select(&entity, &query).unwrap();
1309
1310        let chunks = executor.fetch_stream(&compiled, 10).unwrap();
1311
1312        // Empty result = 1 chunk with 0 rows, marked as last
1313        assert_eq!(chunks.len(), 1);
1314        assert_eq!(chunks[0].rows.len(), 0);
1315        assert!(chunks[0].is_last);
1316    }
1317
1318    #[test]
1319    fn sqlite_fetch_stream_exact_chunk_boundary() {
1320        let executor = SqliteMutationExecutor::new(Arc::new(Mutex::new(
1321            Connection::open_in_memory().unwrap(),
1322        )));
1323        let entity = entity();
1324
1325        executor
1326            .execute(&CompiledQuery {
1327                sql: "CREATE TABLE orders (id INTEGER PRIMARY KEY, version INTEGER, name VARCHAR(255))"
1328                    .to_owned(),
1329                params: Vec::new(),
1330                comment: None,
1331            })
1332            .unwrap();
1333
1334        // Insert exactly 20 rows
1335        for i in 1..=20 {
1336            let insert = SqliteDialect
1337                .compile_insert(
1338                    &entity,
1339                    &InsertCommand::new("Order")
1340                        .value("id", i as u64)
1341                        .value("version", 1_i64)
1342                        .value("name", format!("order-{i}")),
1343                )
1344                .unwrap();
1345            executor.execute(&insert).unwrap();
1346        }
1347
1348        let query = SelectQuery::new("Order")
1349            .filter(Expr::gt("version", 0_i64))
1350            .order_asc("id")
1351            .stream(10);
1352
1353        let compiled = SqliteDialect.compile_select(&entity, &query).unwrap();
1354
1355        let chunks = executor.fetch_stream(&compiled, 10).unwrap();
1356
1357        // 20 rows / 10 per chunk = 2 full chunks + 1 empty final chunk
1358        assert_eq!(chunks.len(), 3);
1359        assert_eq!(chunks[0].rows.len(), 10);
1360        assert!(!chunks[0].is_last);
1361        assert_eq!(chunks[1].rows.len(), 10);
1362        assert!(!chunks[1].is_last);
1363        assert_eq!(chunks[2].rows.len(), 0);
1364        assert!(chunks[2].is_last);
1365    }
1366
1367    #[test]
1368    fn test_parse_sqlite_timestamp() {
1369        let ts1 = parse_sqlite_timestamp("2023-01-01 12:30:45").unwrap();
1370        assert!(matches!(ts1, Value::Timestamp(_)));
1371
1372        let ts2 = parse_sqlite_timestamp("2023-01-01").unwrap();
1373        assert!(matches!(ts2, Value::Timestamp(_)));
1374
1375        let ts3 = parse_sqlite_timestamp("2023-01-01T12:30:45Z").unwrap();
1376        assert!(matches!(ts3, Value::Timestamp(_)));
1377
1378        assert!(parse_sqlite_timestamp("invalid").is_err());
1379    }
1380}