1use alopex_core::kv::KVStore;
2use alopex_core::types::TxnMode;
3use alopex_core::KVTransaction;
4use alopex_sql::catalog::CatalogOverlay;
5use alopex_sql::catalog::TxnCatalogView;
6use alopex_sql::executor::query::execute_query_streaming;
7use alopex_sql::executor::query::iterator::VecIterator;
8use alopex_sql::executor::query::RowIterator;
9use alopex_sql::executor::{
10 build_streaming_pipeline, ColumnInfo, ExecutionResult, Executor, QueryRowIterator, Row,
11};
12use alopex_sql::planner::typed_expr::Projection;
13use alopex_sql::storage::{SqlValue, TxnBridge};
14use alopex_sql::AlopexDialect;
15use alopex_sql::Parser;
16use alopex_sql::Planner;
17use alopex_sql::Statement;
18use alopex_sql::StatementKind;
19
20use crate::Database;
21use crate::Error;
22use crate::Result;
23use crate::SqlResult;
24use crate::Transaction;
25
26pub struct StreamingRows<'a> {
32 columns: Vec<ColumnInfo>,
33 iter: Box<dyn RowIterator + 'a>,
34 projection: Projection,
35 schema: Vec<alopex_sql::catalog::ColumnMetadata>,
36}
37
38impl<'a> StreamingRows<'a> {
39 pub fn columns(&self) -> &[ColumnInfo] {
41 &self.columns
42 }
43
44 pub fn next_row(&mut self) -> Result<Option<Vec<SqlValue>>> {
48 match self.iter.next_row() {
49 Some(result) => {
50 let row = result.map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))?;
51 let projected = self.project_row(&row)?;
52 Ok(Some(projected))
53 }
54 None => Ok(None),
55 }
56 }
57
58 fn project_row(&self, row: &Row) -> Result<Vec<SqlValue>> {
60 match &self.projection {
61 Projection::All(names) => {
62 let mut result = Vec::with_capacity(names.len());
64 for name in names {
65 let idx = self
66 .schema
67 .iter()
68 .position(|c| &c.name == name)
69 .ok_or_else(|| {
70 Error::Sql(alopex_sql::SqlError::Execution {
71 message: format!("column not found: {}", name),
72 code: "ALOPEX-E020",
73 })
74 })?;
75 result.push(row.values.get(idx).cloned().unwrap_or(SqlValue::Null));
76 }
77 Ok(result)
78 }
79 Projection::Columns(cols) => {
80 use alopex_sql::executor::evaluator::{evaluate, EvalContext};
81 let ctx = EvalContext::new(&row.values);
82 let mut result = Vec::with_capacity(cols.len());
83 for col in cols {
84 let value = evaluate(&col.expr, &ctx)
85 .map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))?;
86 result.push(value);
87 }
88 Ok(result)
89 }
90 }
91 }
92}
93
94pub enum StreamingQueryResult<R> {
96 Success,
98 RowsAffected(u64),
100 QueryProcessed(R),
102}
103
104pub enum SqlStreamingResult {
109 Success,
111 RowsAffected(u64),
113 Query(QueryRowIterator<'static>),
115}
116
117fn parse_sql(sql: &str) -> Result<Vec<Statement>> {
118 let dialect = AlopexDialect;
119 Parser::parse_sql(&dialect, sql).map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))
120}
121
122fn stmt_requires_write(stmt: &Statement) -> bool {
123 !matches!(stmt.kind, StatementKind::Select(_))
124}
125
126fn stmt_changes_catalog(stmt: &Statement) -> bool {
127 matches!(
128 stmt.kind,
129 StatementKind::CreateTable(_)
130 | StatementKind::DropTable(_)
131 | StatementKind::CreateIndex(_)
132 | StatementKind::DropIndex(_)
133 )
134}
135
136fn plan_stmt<'a, S: KVStore>(
137 catalog: &'a alopex_sql::catalog::PersistentCatalog<S>,
138 overlay: &'a CatalogOverlay,
139 stmt: &Statement,
140) -> Result<alopex_sql::LogicalPlan> {
141 let view = TxnCatalogView::new(catalog, overlay);
142 let planner = Planner::new(&view);
143 planner
144 .plan(stmt)
145 .map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))
146}
147
148fn build_column_info(
150 projection: &Projection,
151 schema: &[alopex_sql::catalog::ColumnMetadata],
152) -> Result<Vec<ColumnInfo>> {
153 match projection {
154 Projection::All(names) => {
155 let mut cols = Vec::with_capacity(names.len());
156 for name in names {
157 let meta = schema.iter().find(|c| &c.name == name).ok_or_else(|| {
158 Error::Sql(alopex_sql::SqlError::Execution {
159 message: format!("column not found: {}", name),
160 code: "ALOPEX-E020",
161 })
162 })?;
163 cols.push(ColumnInfo::new(name.clone(), meta.data_type.clone()));
164 }
165 Ok(cols)
166 }
167 Projection::Columns(cols) => {
168 let mut result = Vec::with_capacity(cols.len());
169 for (i, col) in cols.iter().enumerate() {
170 let name = col
171 .alias
172 .clone()
173 .or_else(|| {
174 if let alopex_sql::planner::typed_expr::TypedExprKind::ColumnRef {
175 column,
176 ..
177 } = &col.expr.kind
178 {
179 Some(column.clone())
180 } else {
181 None
182 }
183 })
184 .unwrap_or_else(|| format!("col_{}", i));
185 result.push(ColumnInfo::new(name, col.expr.resolved_type.clone()));
186 }
187 Ok(result)
188 }
189 }
190}
191
192impl Database {
193 pub fn execute_sql(&self, sql: &str) -> Result<SqlResult> {
213 Ok(self
214 .execute_sql_multi(sql)?
215 .pop()
216 .unwrap_or(alopex_sql::ExecutionResult::Success))
217 }
218
219 pub fn execute_sql_multi(&self, sql: &str) -> Result<Vec<SqlResult>> {
240 let stmts = parse_sql(sql)?;
241 if stmts.is_empty() {
242 return Ok(Vec::new());
243 }
244
245 if stmts.len() == 1 {
250 let overlay = CatalogOverlay::new();
251 let plan = {
252 let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
253 plan_stmt(&*catalog, &overlay, &stmts[0])?
254 };
255 if matches!(stmts[0].kind, StatementKind::Pragma { .. })
256 || alopex_sql::executor::is_store_direct_plan(&plan)
257 {
258 let mut executor: Executor<_, _> =
259 Executor::new(self.store.clone(), self.sql_catalog.clone());
260 return executor
261 .execute(plan)
262 .map(|result| vec![result])
263 .map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)));
264 }
265 }
266
267 let requires_write = stmts.iter().any(stmt_requires_write);
268 let mode = if requires_write {
269 TxnMode::ReadWrite
270 } else {
271 TxnMode::ReadOnly
272 };
273
274 let mut txn = self.store.begin(mode).map_err(Error::Core)?;
275 let mut overlay = CatalogOverlay::new();
276 let mut borrowed =
277 TxnBridge::<alopex_core::kv::AnyKV>::wrap_external(&mut txn, mode, &mut overlay);
278
279 let mut executor: Executor<_, _> =
280 Executor::new(self.store.clone(), self.sql_catalog.clone());
281
282 let mut results = Vec::with_capacity(stmts.len());
283 for (statement_index, stmt) in stmts.iter().enumerate() {
284 let plan = {
285 let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
286 let (_, overlay) = borrowed.split_parts();
287 plan_stmt(&*catalog, &*overlay, stmt)?
288 };
289
290 {
291 let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
292 let (_, overlay) = borrowed.split_parts();
293 let view = TxnCatalogView::new(&*catalog, &*overlay);
294 self.record_routing(&view, stmt, statement_index);
295 }
296
297 results.push(
298 executor
299 .execute_in_txn(plan, &mut borrowed)
300 .map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))?,
301 );
302 }
303
304 drop(borrowed);
305
306 txn.commit_self().map_err(Error::Core)?;
311 if mode == TxnMode::ReadWrite {
312 let mut catalog = self.sql_catalog.write().expect("catalog lock poisoned");
313 catalog.apply_overlay(overlay);
314 }
315 if stmts.iter().any(stmt_changes_catalog) {
316 self.invalidate_table_info_cache();
317 }
318 if requires_write {
319 let mut cache = self.hnsw_cache.write().expect("hnsw cache lock poisoned");
320 cache.clear();
321 let mut vector_cache = self
322 .vector_cache
323 .write()
324 .expect("vector cache lock poisoned");
325 *vector_cache = None;
326 }
327 Ok(results)
328 }
329
330 pub fn execute_sql_with_rows<F, R>(&self, sql: &str, f: F) -> Result<StreamingQueryResult<R>>
364 where
365 F: FnOnce(StreamingRows<'_>) -> Result<R>,
366 {
367 let stmts = parse_sql(sql)?;
368 if stmts.is_empty() {
369 return Ok(StreamingQueryResult::Success);
370 }
371
372 if stmts.len() == 1 && matches!(stmts[0].kind, StatementKind::Select(_)) {
374 let stmt = &stmts[0];
375 let mode = TxnMode::ReadOnly;
376
377 let mut txn = self.store.begin(mode).map_err(Error::Core)?;
378 let mut overlay = CatalogOverlay::new();
379 let mut borrowed =
380 TxnBridge::<alopex_core::kv::AnyKV>::wrap_external(&mut txn, mode, &mut overlay);
381
382 let plan = {
383 let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
384 let (_, overlay_ref) = borrowed.split_parts();
385 plan_stmt(&*catalog, overlay_ref, stmt)?
386 };
387
388 if alopex_sql::executor::is_store_direct_plan(&plan) {
389 let mut executor: Executor<_, _> =
390 Executor::new(self.store.clone(), self.sql_catalog.clone());
391 let result = executor
392 .execute(plan)
393 .map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))?;
394 let ExecutionResult::Query(query) = result else {
395 return Err(Error::Sql(alopex_sql::SqlError::Execution {
396 message: "store-direct system function did not return rows".into(),
397 code: "ALOPEX-E022",
398 }));
399 };
400 let column_names: Vec<String> = query
401 .columns
402 .iter()
403 .map(|column| column.name.clone())
404 .collect();
405 let schema: Vec<alopex_sql::catalog::ColumnMetadata> = query
406 .columns
407 .iter()
408 .map(|column| {
409 alopex_sql::catalog::ColumnMetadata::new(
410 &column.name,
411 column.data_type.clone(),
412 )
413 })
414 .collect();
415 let rows: Vec<Row> = query
416 .rows
417 .into_iter()
418 .enumerate()
419 .map(|(index, values)| Row::new(index as u64, values))
420 .collect();
421 let iter = VecIterator::new(rows, schema.clone());
422 let streaming_rows = StreamingRows {
423 columns: query.columns,
424 iter: Box::new(iter),
425 projection: Projection::All(column_names),
426 schema,
427 };
428 let result = f(streaming_rows)?;
429 return Ok(StreamingQueryResult::QueryProcessed(result));
430 }
431
432 let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
433 let (mut sql_txn, overlay_ref) = borrowed.split_parts();
434 let view = TxnCatalogView::new(&*catalog, overlay_ref);
435
436 let (iter, projection, schema) = build_streaming_pipeline(&mut sql_txn, &view, plan)
438 .map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))?;
439
440 let columns = build_column_info(&projection, &schema)?;
442
443 let streaming_rows = StreamingRows {
445 columns,
446 iter,
447 projection,
448 schema,
449 };
450
451 let result = f(streaming_rows)?;
453
454 drop(catalog);
456 drop(borrowed);
457 txn.commit_self().map_err(Error::Core)?;
458
459 return Ok(StreamingQueryResult::QueryProcessed(result));
460 }
461
462 let exec_result = self.execute_sql(sql)?;
464 match exec_result {
465 alopex_sql::ExecutionResult::Success => Ok(StreamingQueryResult::Success),
466 alopex_sql::ExecutionResult::RowsAffected(n) => {
467 Ok(StreamingQueryResult::RowsAffected(n))
468 }
469 alopex_sql::ExecutionResult::Query(_qr) => {
470 Err(Error::Sql(alopex_sql::SqlError::Execution {
473 message: "Streaming not available for multi-statement or complex queries"
474 .into(),
475 code: "ALOPEX-E021",
476 }))
477 }
478 }
479 }
480
481 pub fn execute_sql_streaming(&self, sql: &str) -> Result<SqlStreamingResult> {
514 let stmts = parse_sql(sql)?;
515 if stmts.is_empty() {
516 return Ok(SqlStreamingResult::Success);
517 }
518
519 if stmts.len() == 1 && matches!(stmts[0].kind, StatementKind::Select(_)) {
521 let stmt = &stmts[0];
522 let mode = TxnMode::ReadOnly;
523
524 let mut txn = self.store.begin(mode).map_err(Error::Core)?;
525 let mut overlay = CatalogOverlay::new();
526 let mut borrowed =
527 TxnBridge::<alopex_core::kv::AnyKV>::wrap_external(&mut txn, mode, &mut overlay);
528
529 let plan = {
530 let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
531 let (_, overlay) = borrowed.split_parts();
532 plan_stmt(&*catalog, &*overlay, stmt)?
533 };
534
535 let (mut sql_txn, _overlay) = borrowed.split_parts();
536
537 let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
538 let view = TxnCatalogView::new(&*catalog, _overlay);
539 let iter = execute_query_streaming(&mut sql_txn, &view, plan)
540 .map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))?;
541
542 drop(catalog);
543 drop(borrowed);
544
545 txn.commit_self().map_err(Error::Core)?;
546
547 return Ok(SqlStreamingResult::Query(iter));
548 }
549
550 let result = self.execute_sql(sql)?;
552 match result {
553 alopex_sql::ExecutionResult::Success => Ok(SqlStreamingResult::Success),
554 alopex_sql::ExecutionResult::RowsAffected(n) => Ok(SqlStreamingResult::RowsAffected(n)),
555 alopex_sql::ExecutionResult::Query(qr) => {
556 use alopex_sql::executor::query::iterator::VecIterator;
558 use alopex_sql::executor::Row;
559 use alopex_sql::planner::typed_expr::Projection;
560
561 let column_names: Vec<String> = qr.columns.iter().map(|c| c.name.clone()).collect();
562 let schema: Vec<alopex_sql::catalog::ColumnMetadata> = qr
563 .columns
564 .iter()
565 .map(|c| alopex_sql::catalog::ColumnMetadata::new(&c.name, c.data_type.clone()))
566 .collect();
567 let rows: Vec<Row> = qr
568 .rows
569 .into_iter()
570 .enumerate()
571 .map(|(i, values)| Row::new(i as u64, values))
572 .collect();
573 let iter = VecIterator::new(rows, schema.clone());
574 let query_iter =
575 QueryRowIterator::new(Box::new(iter), Projection::All(column_names), schema);
576 Ok(SqlStreamingResult::Query(query_iter))
577 }
578 }
579 }
580}
581
582impl<'a> Transaction<'a> {
583 pub fn execute_sql(&mut self, sql: &str) -> Result<SqlResult> {
600 let stmts = parse_sql(sql)?;
601 if stmts.is_empty() {
602 return Ok(alopex_sql::ExecutionResult::Success);
603 }
604
605 if stmts.iter().any(stmt_requires_write) {
606 let mut cache = self
607 .db
608 .hnsw_cache
609 .write()
610 .expect("hnsw cache lock poisoned");
611 cache.clear();
612 let mut vector_cache = self
613 .db
614 .vector_cache
615 .write()
616 .expect("vector cache lock poisoned");
617 *vector_cache = None;
618 }
619
620 let store = self.db.store.clone();
621 let sql_catalog = self.db.sql_catalog.clone();
622
623 let txn = self.inner.as_mut().ok_or(Error::TxnCompleted)?;
624 let mode = txn.mode();
625
626 let mut borrowed =
627 TxnBridge::<alopex_core::kv::AnyKV>::wrap_external(txn, mode, &mut self.overlay);
628 let mut executor: Executor<_, _> = Executor::new(store, sql_catalog.clone());
629
630 let mut last = alopex_sql::ExecutionResult::Success;
631 for (statement_index, stmt) in stmts.iter().enumerate() {
632 let plan = {
633 let catalog = sql_catalog.read().expect("catalog lock poisoned");
634 let (_, overlay) = borrowed.split_parts();
635 plan_stmt(&*catalog, &*overlay, stmt)?
636 };
637
638 {
639 let catalog = sql_catalog.read().expect("catalog lock poisoned");
640 let (_, overlay) = borrowed.split_parts();
641 let view = TxnCatalogView::new(&*catalog, &*overlay);
642 self.db.record_routing(&view, stmt, statement_index);
643 }
644
645 last = executor
646 .execute_in_txn(plan, &mut borrowed)
647 .map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))?;
648 }
649
650 if stmts.iter().any(stmt_changes_catalog) {
651 self.catalog_modified = true;
652 }
653 Ok(last)
654 }
655}