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 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 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 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 #[allow(clippy::await_holding_lock)]
323 fn stream_sql(
324 &self,
325 query: CompiledQuery,
326 chunk_size: usize,
327 ) -> teaql_data_service::QueryStream<'_, Self::Error> {
328 let connection = self.connection.clone();
329 Box::pin(async_stream::try_stream! {
330 let params = bind_values(&query.params)?;
331 let guard = connection.lock().map_err(|err| MutationExecutorError::Lock(err.to_string()))?;
332 let mut statement = guard.prepare(&query.sql_with_comment())?;
333 let columns = statement_columns(&statement);
334 let mut rows = statement.query(params_from_iter(params.iter()))?;
335 let mut chunk = Vec::with_capacity(chunk_size); let mut index = 0;
336 while let Some(row) = rows.next()? {
337 chunk.push(decode_sqlite_row(row, &columns)?);
338 if chunk.len() == chunk_size { yield teaql_data_service::StreamChunk { rows: std::mem::take(&mut chunk), chunk_index: index, is_last: false }; index += 1; }
339 }
340 if !chunk.is_empty() { yield teaql_data_service::StreamChunk { rows: chunk, chunk_index: index, is_last: true }; }
341 })
342 }
343}
344
345impl teaql_data_service::StreamQueryExecutor for SqliteMutationExecutor {
346 fn query_stream(
347 &self,
348 request: teaql_data_service::QueryRequest,
349 chunk_size: usize,
350 ) -> teaql_data_service::QueryStream<'_, Self::Error> {
351 let dialect = SqliteDialect;
352 let entity_desc = teaql_core::EntityDescriptor::new(&request.query.entity);
354 match dialect.compile_select(&entity_desc, &request.query) {
355 Ok(compiled) => {
356 teaql_sql::StreamingSqlTransport::stream_sql(self, compiled, chunk_size)
357 }
358 Err(error) => Box::pin(futures_util::stream::once(async {
359 Err(MutationExecutorError::SqlCompile(error))
360 })),
361 }
362 }
363}
364
365impl teaql_sql::SqlTransaction for SqliteMutationExecutor {
366 type Error = MutationExecutorError;
367
368 async fn commit_sql(self) -> Result<(), Self::Error> {
369 self.commit_transaction()
370 }
371
372 async fn rollback_sql(self) -> Result<(), Self::Error> {
373 self.rollback_transaction()
374 }
375}
376
377impl teaql_sql::SqlTransactionTransport for SqliteMutationExecutor {
378 type Tx<'a>
379 = Self
380 where
381 Self: 'a;
382
383 async fn begin_sql(&self) -> Result<Self::Tx<'_>, Self::Error> {
384 self.begin_transaction()?;
385 Ok(self.clone())
386 }
387}
388
389fn initial_graph_exists_sqlite(
390 executor: &SqliteMutationExecutor,
391 dialect: &SqliteDialect,
392 entity: &EntityDescriptor,
393 graph: &GraphNode,
394) -> Result<bool, MutationExecutorError> {
395 let Some(id) = graph.values.get("id") else {
396 return Ok(false);
397 };
398 let query = dialect.compile_select(
399 entity,
400 &SelectQuery::new(&graph.entity)
401 .project("id")
402 .filter(Expr::eq("id", id.clone()))
403 .limit(1),
404 )?;
405 Ok(!executor.fetch_all(&query)?.is_empty())
406}
407
408fn compile_initial_graph_insert(
409 dialect: &impl SqlDialect,
410 entity: &EntityDescriptor,
411 graph: &GraphNode,
412) -> Result<CompiledQuery, MutationExecutorError> {
413 let mut command = InsertCommand::new(&graph.entity);
414 for (field, value) in &graph.values {
415 command = command.value(field.clone(), value.clone());
416 }
417 dialect.compile_insert(entity, &command).map_err(Into::into)
418}
419
420fn compile_initial_graph_update(
421 dialect: &impl SqlDialect,
422 entity: &EntityDescriptor,
423 graph: &GraphNode,
424) -> Result<Option<CompiledQuery>, MutationExecutorError> {
425 let Some(id) = graph.values.get("id") else {
426 return Ok(None);
427 };
428 let mut command = UpdateCommand::new(&graph.entity, id.clone());
429 for (field, value) in &graph.values {
430 if field == "id" {
431 continue;
432 }
433 command = command.value(field.clone(), value.clone());
434 }
435 match dialect.compile_update(entity, &command) {
436 Ok(query) => Ok(Some(query)),
437 Err(SqlCompileError::EmptyMutation(_)) => Ok(None),
438 Err(err) => Err(err.into()),
439 }
440}
441
442pub trait SqliteSchemaExt {
443 fn ensure_sqlite_schema(
444 &self,
445 ) -> Pin<Box<dyn Future<Output = Result<(), MutationExecutorError>> + Send + '_>>;
446}
447
448pub fn ensure_sqlite_schema_for(ctx: &UserContext) -> Result<(), MutationExecutorError> {
449 let dialect = ctx.get_resource::<SqliteDialect>().ok_or_else(|| {
450 MutationExecutorError::Bind("missing typed resource: SqliteDialect".to_owned())
451 })?;
452 let executor = ctx
453 .get_resource::<SqliteMutationExecutor>()
454 .ok_or_else(|| {
455 MutationExecutorError::Bind("missing typed resource: SqliteMutationExecutor".to_owned())
456 })?;
457
458 let entities = ctx.all_entities();
459
460 executor.ensure_id_space_table(DEFAULT_ID_SPACE_TABLE)?;
462
463 for entity in &entities {
465 let field_count = entity.properties.len();
466 if !executor.table_exists(&entity.table_name)? {
467 let sql = dialect.compile_create_table(entity)?;
469 executor.lock()?.execute(&sql, [])?;
470 let _ = ctx.send_event(RawAuditEvent::schema_created(
471 &entity.name,
472 &entity.table_name,
473 field_count,
474 ));
475 continue;
476 }
477 let existing_columns = executor.table_columns(&entity.table_name)?;
479 let mut fields_added = 0;
480 for property in &entity.properties {
481 let bare_column = strip_identifier_quotes(&property.column_name).to_lowercase();
482 if existing_columns.contains(&bare_column) {
483 continue;
484 }
485 let sql = dialect.compile_add_column(entity, property)?;
486 executor.lock()?.execute(&sql, [])?;
487 let _ = ctx.send_event(RawAuditEvent::field_added(
488 &entity.name,
489 &entity.table_name,
490 &property.column_name,
491 ));
492 fields_added += 1;
493 }
494 let _ = ctx.send_event(RawAuditEvent::schema_verified(
495 &entity.name,
496 &entity.table_name,
497 field_count,
498 ));
499 let _ = fields_added; }
501
502 let mut seed_counts: BTreeMap<String, (usize, usize)> = BTreeMap::new(); for graph in ctx.initial_graphs() {
505 let entity = ctx.entity(&graph.entity).ok_or_else(|| {
506 MutationExecutorError::Bind(format!("missing entity: {}", graph.entity))
507 })?;
508 let counts = seed_counts.entry(graph.entity.clone()).or_insert((0, 0));
509 if initial_graph_exists_sqlite(executor, dialect, entity, graph)? {
510 if let Some(query) = compile_initial_graph_update(dialect, entity, graph)? {
511 executor.execute(&query)?;
512 }
513 counts.1 += 1; continue;
515 }
516 let query = compile_initial_graph_insert(dialect, entity, graph)?;
517 executor.execute(&query)?;
518 counts.0 += 1; }
520
521 for (entity_name, (inserted, updated)) in &seed_counts {
523 let entity = ctx.entity(entity_name).ok_or_else(|| {
524 MutationExecutorError::Bind(format!("missing entity: {}", entity_name))
525 })?;
526 let _ = ctx.send_event(RawAuditEvent::data_seeded(
527 entity_name,
528 &entity.table_name,
529 *inserted,
530 *updated,
531 ));
532 }
533
534 Ok(())
535}
536
537impl SqliteSchemaExt for UserContext {
538 fn ensure_sqlite_schema(
539 &self,
540 ) -> Pin<Box<dyn Future<Output = Result<(), MutationExecutorError>> + Send + '_>> {
541 Box::pin(async move { ensure_sqlite_schema_for(self) })
542 }
543}
544
545#[derive(Debug, Default, Clone, Copy)]
546pub struct SqliteSchemaProvider;
547
548impl SchemaProvider for SqliteSchemaProvider {
549 fn ensure_schema<'a>(
550 &'a self,
551 ctx: &'a UserContext,
552 ) -> Pin<Box<dyn Future<Output = Result<(), RuntimeError>> + Send + 'a>> {
553 Box::pin(async move {
554 ensure_sqlite_schema_for(ctx).map_err(|err| RuntimeError::Schema(err.to_string()))
555 })
556 }
557}
558
559pub trait SqliteProviderExt {
560 fn use_sqlite_provider(&mut self, executor: SqliteMutationExecutor) -> &mut Self;
561}
562
563impl SqliteProviderExt for UserContext {
564 fn use_sqlite_provider(&mut self, executor: SqliteMutationExecutor) -> &mut Self {
565 self.insert_resource(SqliteDialect);
566 self.insert_resource(executor);
567 self.set_schema_provider(SqliteSchemaProvider);
568 self
569 }
570}
571
572#[derive(Clone)]
573pub struct SqliteIdSpaceGenerator {
574 executor: SqliteMutationExecutor,
575 table_name: String,
576}
577
578impl SqliteIdSpaceGenerator {
579 pub fn new(connection: Connection) -> Self {
580 Self::from_executor(SqliteMutationExecutor::from_connection(connection))
581 }
582
583 pub fn from_executor(executor: SqliteMutationExecutor) -> Self {
584 Self {
585 executor,
586 table_name: DEFAULT_ID_SPACE_TABLE.to_owned(),
587 }
588 }
589
590 pub fn with_table_name(mut self, table_name: impl Into<String>) -> Self {
591 self.table_name = table_name.into();
592 self
593 }
594
595 pub fn ensure_table(&self) -> Result<(), MutationExecutorError> {
596 self.executor.ensure_id_space_table(&self.table_name)
597 }
598
599 pub fn next_id(&self, entity: &str) -> Result<u64, MutationExecutorError> {
600 self.ensure_table()?;
601 let sql = format!(
602 "INSERT INTO {} (type_name, current_level) VALUES (?, 1) \
603 ON CONFLICT (type_name) DO UPDATE \
604 SET current_level = current_level + 1 \
605 RETURNING current_level",
606 quote_ident(&self.table_name)
607 );
608 let id: i64 = self
609 .executor
610 .lock()?
611 .query_row(&sql, [entity], |row| row.get(0))?;
612 u64::try_from(id).map_err(|_| {
613 MutationExecutorError::Bind(format!("generated id {id} cannot be represented as u64"))
614 })
615 }
616}
617
618impl InternalIdGenerator for SqliteIdSpaceGenerator {
619 fn generate_id(&self, entity: &str) -> Result<u64, RuntimeError> {
620 self.next_id(entity)
621 .map_err(|err| RuntimeError::IdGeneration(err.to_string()))
622 }
623}
624
625fn quote_ident(ident: &str) -> String {
626 quote_identifier_if_needed(ident, '"')
627}
628
629fn strip_identifier_quotes(ident: &str) -> &str {
637 let bytes = ident.as_bytes();
638 if bytes.len() >= 2 {
639 let (first, last) = (bytes[0], bytes[bytes.len() - 1]);
640 if (first == b'"' && last == b'"')
641 || (first == b'`' && last == b'`')
642 || (first == b'[' && last == b']')
643 {
644 return &ident[1..ident.len() - 1];
645 }
646 }
647 ident
648}
649
650fn bind_values(values: &[Value]) -> Result<Vec<SqliteValue>, MutationExecutorError> {
651 values.iter().map(bind_sqlite_value).collect()
652}
653
654fn bind_sqlite_value(value: &Value) -> Result<SqliteValue, MutationExecutorError> {
655 match value {
656 Value::Null => Ok(SqliteValue::Null),
657 Value::Bool(v) => Ok(SqliteValue::Integer(i64::from(*v))),
658 Value::I64(v) => Ok(SqliteValue::Integer(*v)),
659 Value::U64(v) => i64::try_from(*v)
660 .map(SqliteValue::Integer)
661 .map_err(|_| MutationExecutorError::Bind(format!("u64 value {v} exceeds i64 range"))),
662 Value::F64(v) => Ok(SqliteValue::Real(*v)),
663 Value::Decimal(v) => Ok(SqliteValue::Text(v.to_string())),
667 Value::Text(v) => Ok(SqliteValue::Text(v.clone())),
668 Value::Json(v) => Ok(SqliteValue::Text(v.to_string())),
669 Value::Date(v) => Ok(SqliteValue::Text(v.format("%Y-%m-%d").to_string())),
670 Value::Timestamp(v) => Ok(SqliteValue::Text(v.0.to_string())),
671 Value::Object(_) => Err(MutationExecutorError::UnsupportedValue("object")),
672 Value::List(_) => Err(MutationExecutorError::UnsupportedValue("list")),
673 Value::TypedNull(_) => Ok(SqliteValue::Null),
674 }
675}
676
677#[derive(Debug, Clone)]
678struct ColumnInfo {
679 name: String,
680 decl_type: Option<String>,
681}
682
683fn statement_columns(statement: &rusqlite::Statement<'_>) -> Vec<ColumnInfo> {
684 statement
685 .columns()
686 .into_iter()
687 .map(|column| ColumnInfo {
688 name: column.name().to_owned(),
689 decl_type: column.decl_type().map(|value| value.to_ascii_uppercase()),
690 })
691 .collect()
692}
693
694fn decode_sqlite_row(
695 row: &Row<'_>,
696 columns: &[ColumnInfo],
697) -> Result<Record, MutationExecutorError> {
698 let mut record = BTreeMap::new();
699 for (index, column) in columns.iter().enumerate() {
700 let value_ref = row.get_ref(index)?;
701 let value = match value_ref {
702 ValueRef::Null => Value::Null,
703 ValueRef::Integer(value) => decode_sqlite_integer(value, column),
704 ValueRef::Real(value) => Value::F64(value),
705 ValueRef::Text(value) => decode_sqlite_text(value, column)?,
706 ValueRef::Blob(_) => {
707 return Err(MutationExecutorError::UnsupportedColumnType(
708 "BLOB".to_owned(),
709 ));
710 }
711 };
712 record.insert(column.name.clone(), value);
713 }
714 Ok(record)
715}
716
717fn decode_sqlite_integer(value: i64, column: &ColumnInfo) -> Value {
718 match column_decl_type(column).as_deref() {
719 Some("BOOLEAN") | Some("BOOL") => Value::Bool(value != 0),
720 _ => Value::I64(value),
721 }
722}
723
724fn decode_sqlite_text(value: &[u8], column: &ColumnInfo) -> Result<Value, MutationExecutorError> {
725 let value = std::str::from_utf8(value)
726 .map_err(|err| MutationExecutorError::Bind(format!("invalid sqlite text: {err}")))?;
727 match column_decl_type(column).as_deref() {
728 Some("NUMERIC") | Some("DECIMAL") => Decimal::from_str(value)
729 .map(Value::Decimal)
730 .map_err(|err| MutationExecutorError::Bind(format!("invalid sqlite decimal: {err}"))),
731 Some("JSON") => serde_json::from_str(value).map(Value::Json).map_err(|err| {
732 MutationExecutorError::Bind(format!("invalid sqlite json value: {err}"))
733 }),
734 Some("DATE") => NaiveDate::parse_from_str(value, "%Y-%m-%d")
735 .map(Value::Date)
736 .map_err(|err| MutationExecutorError::Bind(format!("invalid sqlite date: {err}"))),
737 Some("TIMESTAMP") | Some("DATETIME") => parse_sqlite_timestamp(value),
738 _ => infer_sqlite_text(value),
739 }
740}
741
742fn infer_sqlite_text(value: &str) -> Result<Value, MutationExecutorError> {
743 if let Ok(date) = NaiveDate::parse_from_str(value, "%Y-%m-%d") {
744 return Ok(Value::Date(date));
745 }
746 if let Ok(timestamp) = DateTime::parse_from_rfc3339(value) {
747 return Ok(Value::Timestamp(teaql_core::time::Timestamp(
748 timestamp.timestamp_millis(),
749 )));
750 }
751 if let Ok(timestamp) = NaiveDateTime::parse_from_str(value, "%Y-%m-%d %H:%M:%S") {
752 return Ok(Value::Timestamp(teaql_core::time::Timestamp(
753 timestamp.and_utc().timestamp_millis(),
754 )));
755 }
756 Ok(Value::Text(value.to_owned()))
757}
758
759fn parse_sqlite_timestamp(value: &str) -> Result<Value, MutationExecutorError> {
760 if let Ok(timestamp) = DateTime::parse_from_rfc3339(value) {
761 return Ok(Value::Timestamp(teaql_core::time::Timestamp(
762 timestamp.timestamp_millis(),
763 )));
764 }
765 if let Ok(date) = NaiveDate::parse_from_str(value, "%Y-%m-%d") {
766 return Ok(Value::Timestamp(teaql_core::time::Timestamp(
767 date.and_hms_opt(0, 0, 0)
768 .unwrap_or_default()
769 .and_utc()
770 .timestamp_millis(),
771 )));
772 }
773 NaiveDateTime::parse_from_str(value, "%Y-%m-%d %H:%M:%S")
774 .map(|timestamp| {
775 Value::Timestamp(teaql_core::time::Timestamp(
776 timestamp.and_utc().timestamp_millis(),
777 ))
778 })
779 .map_err(|err| MutationExecutorError::Bind(format!("invalid sqlite timestamp: {err}")))
780}
781
782fn column_decl_type(column: &ColumnInfo) -> Option<String> {
783 column
784 .decl_type
785 .as_ref()
786 .map(|value| value.split('(').next().unwrap_or(value).trim().to_owned())
787}
788
789#[cfg(test)]
790mod tests {
791 use super::*;
792 use futures_util::StreamExt;
793 use teaql_core::{DeleteCommand, RecoverCommand};
794 use teaql_macros::TeaqlEntity;
795 use teaql_runtime::InMemoryMetadataStore;
796
797 #[test]
798 fn streaming_sql_yields_bounded_chunks_and_releases_cursor_on_drop() {
799 let connection = Connection::open_in_memory().unwrap();
800 connection
801 .execute_batch(
802 "CREATE TABLE stream_fixture(id INTEGER);\
803 INSERT INTO stream_fixture VALUES (1), (2), (3), (4), (5);",
804 )
805 .unwrap();
806 let executor = SqliteMutationExecutor::from_connection(connection);
807 let query = CompiledQuery {
808 sql: "SELECT id FROM stream_fixture ORDER BY id".to_owned(),
809 params: vec![],
810 comment: None,
811 };
812 let mut stream = teaql_sql::StreamingSqlTransport::stream_sql(&executor, query.clone(), 2);
813 let sizes = futures_executor::block_on(async {
814 let mut result = Vec::new();
815 while let Some(chunk) = stream.next().await {
816 result.push(chunk.unwrap().rows.len());
817 }
818 result
819 });
820 assert_eq!(sizes, vec![2, 2, 1]);
821
822 let mut early = teaql_sql::StreamingSqlTransport::stream_sql(&executor, query, 2);
823 assert_eq!(
824 futures_executor::block_on(early.next())
825 .unwrap()
826 .unwrap()
827 .rows
828 .len(),
829 2
830 );
831 drop(early);
832 let count: i64 = executor
833 .connection()
834 .lock()
835 .unwrap()
836 .query_row("SELECT count(*) FROM stream_fixture", [], |row| row.get(0))
837 .unwrap();
838 assert_eq!(count, 5);
839 }
840
841 #[test]
842 fn decimal_bind_is_numeric_and_comparable() {
843 let value =
844 bind_sqlite_value(&Value::Decimal(Decimal::from_str("123.450").unwrap())).unwrap();
845 assert_eq!(value, SqliteValue::Text("123.450".to_owned()));
846 let connection = Connection::open_in_memory().unwrap();
847 let matches: i64 = connection
848 .query_row(
849 "SELECT 1 WHERE CAST(? AS NUMERIC) BETWEEN 120 AND 130",
850 [value],
851 |row| row.get(0),
852 )
853 .unwrap();
854 assert_eq!(matches, 1);
855 }
856
857 fn entity() -> EntityDescriptor {
858 EntityDescriptor::new("Order")
859 .table_name("orders")
860 .property(
861 PropertyDescriptor::new("id", DataType::U64)
862 .column_name("id")
863 .id()
864 .not_null(),
865 )
866 .property(
867 PropertyDescriptor::new("version", DataType::I64)
868 .column_name("version")
869 .version()
870 .not_null(),
871 )
872 .property(PropertyDescriptor::new("name", DataType::Text).column_name("name"))
873 }
874
875 fn order_line_entity() -> EntityDescriptor {
876 EntityDescriptor::new("OrderLine")
877 .table_name("order_line")
878 .property(
879 PropertyDescriptor::new("id", DataType::U64)
880 .column_name("id")
881 .id()
882 .not_null(),
883 )
884 .property(
885 PropertyDescriptor::new("order_id", DataType::U64)
886 .column_name("order_id")
887 .not_null(),
888 )
889 .property(PropertyDescriptor::new("name", DataType::Text).column_name("name"))
890 }
891
892 #[allow(dead_code)]
893 #[derive(Debug, PartialEq, TeaqlEntity)]
894 #[teaql(entity = "FeatureFlag", table = "feature_flags")]
895 struct FeatureFlagRow {
896 #[teaql(id)]
897 id: u64,
898 #[teaql(version)]
899 version: i64,
900 enabled: bool,
901 optional_enabled: Option<bool>,
902 }
903
904 fn feature_flag_record(enabled: Value, optional_enabled: Value) -> Record {
905 Record::from([
906 ("id".to_owned(), Value::U64(1)),
907 ("version".to_owned(), Value::I64(1)),
908 ("enabled".to_owned(), enabled),
909 ("optional_enabled".to_owned(), optional_enabled),
910 ])
911 }
912
913 #[test]
914 fn sqlite_dialect_compiles_mutations_and_schema() {
915 let insert = SqliteDialect
916 .compile_insert(
917 &entity(),
918 &InsertCommand::new("Order")
919 .value("id", 1_u64)
920 .value("name", "A"),
921 )
922 .unwrap();
923 assert_eq!(insert.sql, "INSERT INTO orders (id, name) VALUES (?, ?)");
924
925 let update = SqliteDialect
926 .compile_update(
927 &entity(),
928 &UpdateCommand::new("Order", 1_u64)
929 .expected_version(3)
930 .value("name", "B"),
931 )
932 .unwrap();
933 assert_eq!(
934 update.sql,
935 "UPDATE orders SET name = ?, version = ? WHERE id = ? AND version = ?"
936 );
937
938 let delete = SqliteDialect
939 .compile_delete(
940 &entity(),
941 &DeleteCommand::new("Order", 1_u64).expected_version(3),
942 )
943 .unwrap();
944 let recover = SqliteDialect
945 .compile_recover(&entity(), &RecoverCommand::new("Order", 1_u64, -4))
946 .unwrap();
947 assert_eq!(
948 delete.sql,
949 "UPDATE orders SET version = ? WHERE id = ? AND version = ?"
950 );
951 assert_eq!(
952 recover.sql,
953 "UPDATE orders SET version = ? WHERE id = ? AND version = ?"
954 );
955
956 let create = SqliteDialect.compile_create_table(&entity()).unwrap();
957 assert_eq!(
958 create,
959 "CREATE TABLE IF NOT EXISTS orders (id INTEGER PRIMARY KEY NOT NULL, version INTEGER NOT NULL, name VARCHAR(255))"
960 );
961 }
962
963 #[test]
964 fn sqlite_executor_ensures_schema_and_roundtrips_rows() {
965 let executor =
966 SqliteMutationExecutor::from_connection(Connection::open_in_memory().unwrap());
967 let entity = entity();
968 let mut ctx = UserContext::new()
969 .with_metadata(InMemoryMetadataStore::new().with_entity(entity.clone()));
970
971 ctx.use_sqlite_provider(executor.clone());
972 ensure_sqlite_schema_for(&ctx).unwrap();
973
974 let insert = SqliteDialect
975 .compile_insert(
976 &entity,
977 &InsertCommand::new("Order")
978 .value("id", 1_u64)
979 .value("version", 1_i64)
980 .value("name", "draft"),
981 )
982 .unwrap();
983 assert_eq!(executor.execute(&insert).unwrap(), 1);
984
985 let select = SqliteDialect
986 .compile_select(
987 &entity,
988 &SelectQuery::new("Order")
989 .filter(Expr::eq("id", 1_u64))
990 .order_asc("id"),
991 )
992 .unwrap();
993 let rows = executor.fetch_all(&select).unwrap();
994 assert_eq!(rows.len(), 1);
995 assert_eq!(rows[0].get("id"), Some(&Value::I64(1)));
996 assert_eq!(rows[0].get("version"), Some(&Value::I64(1)));
997 assert_eq!(rows[0].get("name"), Some(&Value::Text("draft".to_owned())));
998 }
999
1000 #[test]
1001 fn sqlite_executes_partitioned_relation_limit_per_parent() {
1002 let executor =
1003 SqliteMutationExecutor::from_connection(Connection::open_in_memory().unwrap());
1004 let entity = order_line_entity();
1005 executor.ensure_schema(&SqliteDialect, &[&entity]).unwrap();
1006
1007 for order_id in [11_u64, 12_u64] {
1008 for index in 1_u64..=5 {
1009 let id = order_id * 100 + index;
1010 let insert = SqliteDialect
1011 .compile_insert(
1012 &entity,
1013 &InsertCommand::new("OrderLine")
1014 .value("id", id)
1015 .value("order_id", order_id)
1016 .value("name", format!("line-{id}")),
1017 )
1018 .unwrap();
1019 executor.execute(&insert).unwrap();
1020 }
1021 }
1022
1023 let query = SelectQuery::new("OrderLine")
1024 .project("id")
1025 .project("order_id")
1026 .order_desc("id")
1027 .limit(3)
1028 .partition_by("order_id");
1029 let compiled = SqliteDialect.compile_select(&entity, &query).unwrap();
1030 let rows = executor.fetch_all(&compiled).unwrap();
1031
1032 assert_eq!(rows.len(), 6);
1033 for order_id in [11_i64, 12_i64] {
1034 let ids = rows
1035 .iter()
1036 .filter(|row| row.get("order_id") == Some(&Value::I64(order_id)))
1037 .filter_map(|row| row.get("id").cloned())
1038 .collect::<Vec<_>>();
1039 assert_eq!(
1040 ids,
1041 vec![
1042 Value::I64(order_id * 100 + 5),
1043 Value::I64(order_id * 100 + 4),
1044 Value::I64(order_id * 100 + 3),
1045 ]
1046 );
1047 }
1048 }
1049
1050 #[test]
1051 fn sqlite_boolean_new_schema_roundtrips_as_bool() {
1052 let executor =
1053 SqliteMutationExecutor::from_connection(Connection::open_in_memory().unwrap());
1054 let entity = <FeatureFlagRow as teaql_core::TeaqlEntity>::entity_descriptor();
1055 let ddl = SqliteDialect.compile_create_table(&entity).unwrap();
1056 assert!(ddl.contains("enabled BOOLEAN NOT NULL"), "{ddl}");
1057 assert!(ddl.contains("optional_enabled BOOLEAN"), "{ddl}");
1058 assert!(!ddl.contains("enabled INTEGER"), "{ddl}");
1059
1060 executor.ensure_schema(&SqliteDialect, &[&entity]).unwrap();
1061 for (id, enabled, optional_enabled) in [(1_u64, false, true), (2_u64, true, false)] {
1062 let insert = SqliteDialect
1063 .compile_insert(
1064 &entity,
1065 &InsertCommand::new("FeatureFlag")
1066 .value("id", id)
1067 .value("version", 1_i64)
1068 .value("enabled", enabled)
1069 .value("optional_enabled", optional_enabled),
1070 )
1071 .unwrap();
1072 assert_eq!(executor.execute(&insert).unwrap(), 1);
1073 }
1074
1075 let select = SqliteDialect
1076 .compile_select(&entity, &SelectQuery::new("FeatureFlag").order_asc("id"))
1077 .unwrap();
1078 let rows = executor.fetch_all(&select).unwrap();
1079 assert_eq!(rows[0].get("enabled"), Some(&Value::Bool(false)));
1080 assert_eq!(rows[0].get("optional_enabled"), Some(&Value::Bool(true)));
1081 assert_eq!(rows[1].get("enabled"), Some(&Value::Bool(true)));
1082 assert_eq!(rows[1].get("optional_enabled"), Some(&Value::Bool(false)));
1083
1084 let first = <FeatureFlagRow as teaql_core::Entity>::from_record(rows[0].clone()).unwrap();
1085 let second = <FeatureFlagRow as teaql_core::Entity>::from_record(rows[1].clone()).unwrap();
1086 assert!(!first.enabled);
1087 assert_eq!(first.optional_enabled, Some(true));
1088 assert!(second.enabled);
1089 assert_eq!(second.optional_enabled, Some(false));
1090 }
1091
1092 #[test]
1093 fn sqlite_boolean_legacy_integer_schema_maps_only_binary_values() {
1094 let executor =
1095 SqliteMutationExecutor::from_connection(Connection::open_in_memory().unwrap());
1096 let entity = <FeatureFlagRow as teaql_core::TeaqlEntity>::entity_descriptor();
1097 executor
1098 .execute(&CompiledQuery {
1099 sql: "CREATE TABLE feature_flags (id INTEGER PRIMARY KEY, version INTEGER NOT NULL, enabled INTEGER NOT NULL, optional_enabled INTEGER)"
1100 .to_owned(),
1101 params: Vec::new(),
1102 comment: None,
1103 })
1104 .unwrap();
1105
1106 let insert = SqliteDialect
1107 .compile_insert(
1108 &entity,
1109 &InsertCommand::new("FeatureFlag")
1110 .value("id", 1_u64)
1111 .value("version", 1_i64)
1112 .value("enabled", true)
1113 .value("optional_enabled", false),
1114 )
1115 .unwrap();
1116 executor.execute(&insert).unwrap();
1117 executor
1118 .execute(&CompiledQuery {
1119 sql: "INSERT INTO feature_flags (id, version, enabled, optional_enabled) VALUES (?, ?, ?, ?)"
1120 .to_owned(),
1121 params: vec![
1122 Value::U64(2),
1123 Value::I64(1),
1124 Value::I64(2),
1125 Value::Null,
1126 ],
1127 comment: None,
1128 })
1129 .unwrap();
1130 let select = SqliteDialect
1131 .compile_select(&entity, &SelectQuery::new("FeatureFlag").order_asc("id"))
1132 .unwrap();
1133 let rows = executor.fetch_all(&select).unwrap();
1134 assert_eq!(rows[0].get("version"), Some(&Value::I64(1)));
1135 assert_eq!(rows[0].get("enabled"), Some(&Value::I64(1)));
1136 assert_eq!(rows[0].get("optional_enabled"), Some(&Value::I64(0)));
1137
1138 let decoded = <FeatureFlagRow as teaql_core::Entity>::from_record(rows[0].clone()).unwrap();
1139 assert!(decoded.enabled);
1140 assert_eq!(decoded.optional_enabled, Some(false));
1141 assert_eq!(rows[1].get("enabled"), Some(&Value::I64(2)));
1142 let error =
1143 <FeatureFlagRow as teaql_core::Entity>::from_record(rows[1].clone()).unwrap_err();
1144 assert!(error.message.contains("invalid field enabled"));
1145
1146 for (value, expected) in [
1147 (Value::I64(0), false),
1148 (Value::I64(1), true),
1149 (Value::U64(0), false),
1150 (Value::U64(1), true),
1151 ] {
1152 let decoded = <FeatureFlagRow as teaql_core::Entity>::from_record(feature_flag_record(
1153 value,
1154 Value::Null,
1155 ))
1156 .unwrap();
1157 assert_eq!(decoded.enabled, expected);
1158 assert_eq!(decoded.optional_enabled, None);
1159 }
1160
1161 for invalid in [Value::I64(-1), Value::I64(2), Value::U64(2)] {
1162 let error = <FeatureFlagRow as teaql_core::Entity>::from_record(feature_flag_record(
1163 invalid,
1164 Value::Null,
1165 ))
1166 .unwrap_err();
1167 assert!(error.message.contains("invalid field enabled"));
1168 }
1169 let error = <FeatureFlagRow as teaql_core::Entity>::from_record(feature_flag_record(
1170 Value::Bool(true),
1171 Value::U64(2),
1172 ))
1173 .unwrap_err();
1174 assert!(error.message.contains("invalid field optional_enabled"));
1175 }
1176
1177 #[test]
1178 fn sqlite_executor_parses_json_only_for_json_columns() {
1179 let executor =
1180 SqliteMutationExecutor::from_connection(Connection::open_in_memory().unwrap());
1181
1182 executor
1183 .execute(&CompiledQuery {
1184 sql: "CREATE TABLE payloads (text_payload TEXT, json_payload JSON)".to_owned(),
1185 params: Vec::new(),
1186 comment: None,
1187 })
1188 .unwrap();
1189 executor
1190 .execute(&CompiledQuery {
1191 sql: "INSERT INTO payloads (text_payload, json_payload) VALUES (?, ?)".to_owned(),
1192 params: vec![
1193 Value::Text("{\"active\":true}".to_owned()),
1194 Value::Json(serde_json::json!({"active": true})),
1195 ],
1196 comment: None,
1197 })
1198 .unwrap();
1199
1200 let rows = executor
1201 .fetch_all(&CompiledQuery {
1202 sql: "SELECT text_payload, json_payload FROM payloads".to_owned(),
1203 params: Vec::new(),
1204 comment: None,
1205 })
1206 .unwrap();
1207
1208 assert_eq!(
1209 rows[0].get("text_payload"),
1210 Some(&Value::Text("{\"active\":true}".to_owned()))
1211 );
1212 assert_eq!(
1213 rows[0].get("json_payload"),
1214 Some(&Value::Json(serde_json::json!({"active": true})))
1215 );
1216 }
1217
1218 #[test]
1219 fn sqlite_id_space_generator_increments_ids() {
1220 let executor =
1221 SqliteMutationExecutor::from_connection(Connection::open_in_memory().unwrap());
1222 let generator = SqliteIdSpaceGenerator::from_executor(executor);
1223 assert_eq!(generator.next_id("Order").unwrap(), 1);
1224 assert_eq!(generator.next_id("Order").unwrap(), 2);
1225 }
1226
1227 #[test]
1228 fn sqlite_fetch_stream_returns_chunked_rows() {
1229 let executor = SqliteMutationExecutor::new(Arc::new(Mutex::new(
1230 Connection::open_in_memory().unwrap(),
1231 )));
1232 let entity = entity();
1233
1234 executor
1236 .execute(&CompiledQuery {
1237 sql: "CREATE TABLE orders (id INTEGER PRIMARY KEY, version INTEGER, name VARCHAR(255))"
1238 .to_owned(),
1239 params: Vec::new(),
1240 comment: None,
1241 })
1242 .unwrap();
1243
1244 for i in 1..=25 {
1245 let insert = SqliteDialect
1246 .compile_insert(
1247 &entity,
1248 &InsertCommand::new("Order")
1249 .value("id", i as u64)
1250 .value("version", 1_i64)
1251 .value("name", format!("order-{i}")),
1252 )
1253 .unwrap();
1254 executor.execute(&insert).unwrap();
1255 }
1256
1257 let query = SelectQuery::new("Order")
1259 .filter(Expr::gt("version", 0_i64))
1260 .order_asc("id")
1261 .stream(10);
1262
1263 let compiled = SqliteDialect.compile_select(&entity, &query).unwrap();
1264
1265 let chunks = executor.fetch_stream(&compiled, 10).unwrap();
1266
1267 assert_eq!(chunks.len(), 3);
1269 assert_eq!(chunks[0].rows.len(), 10);
1270 assert_eq!(chunks[0].chunk_index, 0);
1271 assert!(!chunks[0].is_last);
1272
1273 assert_eq!(chunks[1].rows.len(), 10);
1274 assert_eq!(chunks[1].chunk_index, 1);
1275 assert!(!chunks[1].is_last);
1276
1277 assert_eq!(chunks[2].rows.len(), 5);
1278 assert_eq!(chunks[2].chunk_index, 2);
1279 assert!(chunks[2].is_last);
1280
1281 assert_eq!(
1283 chunks[0].rows[0].get("name"),
1284 Some(&Value::Text("order-1".to_owned()))
1285 );
1286 assert_eq!(
1287 chunks[2].rows[4].get("name"),
1288 Some(&Value::Text("order-25".to_owned()))
1289 );
1290 }
1291
1292 #[test]
1293 fn sqlite_fetch_stream_handles_empty_result() {
1294 let executor = SqliteMutationExecutor::new(Arc::new(Mutex::new(
1295 Connection::open_in_memory().unwrap(),
1296 )));
1297
1298 executor
1299 .execute(&CompiledQuery {
1300 sql: "CREATE TABLE orders (id INTEGER PRIMARY KEY, version INTEGER, name VARCHAR(255))"
1301 .to_owned(),
1302 params: Vec::new(),
1303 comment: None,
1304 })
1305 .unwrap();
1306
1307 let entity = entity();
1308 let query = SelectQuery::new("Order")
1309 .filter(Expr::gt("version", 0_i64))
1310 .stream(10);
1311
1312 let compiled = SqliteDialect.compile_select(&entity, &query).unwrap();
1313
1314 let chunks = executor.fetch_stream(&compiled, 10).unwrap();
1315
1316 assert_eq!(chunks.len(), 1);
1318 assert_eq!(chunks[0].rows.len(), 0);
1319 assert!(chunks[0].is_last);
1320 }
1321
1322 #[test]
1323 fn sqlite_fetch_stream_exact_chunk_boundary() {
1324 let executor = SqliteMutationExecutor::new(Arc::new(Mutex::new(
1325 Connection::open_in_memory().unwrap(),
1326 )));
1327 let entity = entity();
1328
1329 executor
1330 .execute(&CompiledQuery {
1331 sql: "CREATE TABLE orders (id INTEGER PRIMARY KEY, version INTEGER, name VARCHAR(255))"
1332 .to_owned(),
1333 params: Vec::new(),
1334 comment: None,
1335 })
1336 .unwrap();
1337
1338 for i in 1..=20 {
1340 let insert = SqliteDialect
1341 .compile_insert(
1342 &entity,
1343 &InsertCommand::new("Order")
1344 .value("id", i as u64)
1345 .value("version", 1_i64)
1346 .value("name", format!("order-{i}")),
1347 )
1348 .unwrap();
1349 executor.execute(&insert).unwrap();
1350 }
1351
1352 let query = SelectQuery::new("Order")
1353 .filter(Expr::gt("version", 0_i64))
1354 .order_asc("id")
1355 .stream(10);
1356
1357 let compiled = SqliteDialect.compile_select(&entity, &query).unwrap();
1358
1359 let chunks = executor.fetch_stream(&compiled, 10).unwrap();
1360
1361 assert_eq!(chunks.len(), 3);
1363 assert_eq!(chunks[0].rows.len(), 10);
1364 assert!(!chunks[0].is_last);
1365 assert_eq!(chunks[1].rows.len(), 10);
1366 assert!(!chunks[1].is_last);
1367 assert_eq!(chunks[2].rows.len(), 0);
1368 assert!(chunks[2].is_last);
1369 }
1370
1371 #[test]
1372 fn test_parse_sqlite_timestamp() {
1373 let ts1 = parse_sqlite_timestamp("2023-01-01 12:30:45").unwrap();
1374 assert!(matches!(ts1, Value::Timestamp(_)));
1375
1376 let ts2 = parse_sqlite_timestamp("2023-01-01").unwrap();
1377 assert!(matches!(ts2, Value::Timestamp(_)));
1378
1379 let ts3 = parse_sqlite_timestamp("2023-01-01T12:30:45Z").unwrap();
1380 assert!(matches!(ts3, Value::Timestamp(_)));
1381
1382 assert!(parse_sqlite_timestamp("invalid").is_err());
1383 }
1384}