use alopex_core::kv::RangeChangeJournalCapability;
use alopex_core::kv::{any::AnyKVTransaction, KVStore, OwnedKVTransactionAdapter, ReadAtPoint};
use alopex_core::types::TxnMode;
use alopex_core::KVTransaction;
use alopex_sql::catalog::TxnCatalogView;
use alopex_sql::catalog::{Catalog, CatalogOverlay};
use alopex_sql::executor::query::execute_query_streaming;
use alopex_sql::executor::query::iterator::VecIterator;
use alopex_sql::executor::query::RowIterator;
use alopex_sql::executor::{
build_streaming_pipeline, ColumnInfo, ExecutionResult, Executor, QueryRowIterator, Row,
};
use alopex_sql::planner::typed_expr::Projection;
use alopex_sql::storage::{LocalRangeChangeJournal, RangeChangeJournalScope, SqlValue, TxnBridge};
use alopex_sql::AlopexDialect;
use alopex_sql::Parser;
use alopex_sql::Planner;
use alopex_sql::Statement;
use alopex_sql::StatementKind;
use std::collections::BTreeMap;
use std::sync::Arc;
use crate::Database;
use crate::Error;
use crate::OwnedEmbeddedTransaction;
use crate::Result;
use crate::SqlResult;
use crate::Transaction;
pub struct StreamingRows<'a> {
columns: Vec<ColumnInfo>,
iter: Box<dyn RowIterator + 'a>,
projection: Projection,
schema: Vec<alopex_sql::catalog::ColumnMetadata>,
}
impl<'a> StreamingRows<'a> {
pub fn columns(&self) -> &[ColumnInfo] {
&self.columns
}
pub fn next_row(&mut self) -> Result<Option<Vec<SqlValue>>> {
match self.iter.next_row() {
Some(result) => {
let row = result.map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))?;
let projected = self.project_row(&row)?;
Ok(Some(projected))
}
None => Ok(None),
}
}
fn project_row(&self, row: &Row) -> Result<Vec<SqlValue>> {
match &self.projection {
Projection::All(names) => {
let mut result = Vec::with_capacity(names.len());
for name in names {
let idx = self
.schema
.iter()
.position(|c| &c.name == name)
.ok_or_else(|| {
Error::Sql(alopex_sql::SqlError::Execution {
message: format!("column not found: {}", name),
code: "ALOPEX-E020",
})
})?;
result.push(row.values.get(idx).cloned().unwrap_or(SqlValue::Null));
}
Ok(result)
}
Projection::Columns(cols) => {
use alopex_sql::executor::evaluator::{evaluate, EvalContext};
let ctx = EvalContext::new(&row.values);
let mut result = Vec::with_capacity(cols.len());
for col in cols {
let value = evaluate(&col.expr, &ctx)
.map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))?;
result.push(value);
}
Ok(result)
}
}
}
}
pub enum StreamingQueryResult<R> {
Success,
RowsAffected(u64),
QueryProcessed(R),
}
pub enum SqlStreamingResult {
Success,
RowsAffected(u64),
Query(QueryRowIterator<'static>),
}
fn parse_sql(sql: &str) -> Result<Vec<Statement>> {
let dialect = AlopexDialect;
Parser::parse_sql(&dialect, sql).map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))
}
fn stmt_requires_write(stmt: &Statement) -> bool {
!matches!(stmt.kind, StatementKind::Select(_))
}
fn stmt_changes_catalog(stmt: &Statement) -> bool {
matches!(
stmt.kind,
StatementKind::CreateTable(_)
| StatementKind::DropTable(_)
| StatementKind::CreateIndex(_)
| StatementKind::DropIndex(_)
)
}
fn stmt_changes_user_data(stmt: &Statement) -> bool {
matches!(
stmt.kind,
StatementKind::Insert(_) | StatementKind::Update(_) | StatementKind::Delete(_)
)
}
pub(crate) fn local_journal_scope<C: Catalog>(catalog: &C) -> RangeChangeJournalScope {
let mut index_tables = BTreeMap::new();
for table in catalog.list_tables() {
for index in catalog.get_indexes_for_table(&table.name) {
index_tables.insert(index.index_id, table.table_id);
}
}
RangeChangeJournalScope::local(index_tables)
}
fn plan_stmt<'a, S: KVStore>(
catalog: &'a alopex_sql::catalog::PersistentCatalog<S>,
overlay: &'a CatalogOverlay,
stmt: &Statement,
) -> Result<alopex_sql::LogicalPlan> {
let view = TxnCatalogView::new(catalog, overlay);
let planner = Planner::new(&view);
planner
.plan(stmt)
.map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))
}
pub(crate) fn execute_sql_owned(
transaction: &mut OwnedEmbeddedTransaction,
sql: &str,
) -> Result<SqlResult> {
let statements = parse_sql(sql)?;
if statements.is_empty() {
return Ok(alopex_sql::ExecutionResult::Success);
}
if statements.iter().any(stmt_requires_write) {
let mut cache = transaction
.db
.hnsw_cache
.write()
.expect("hnsw cache lock poisoned");
cache.clear();
let mut vector_cache = transaction
.db
.vector_cache
.write()
.expect("vector cache lock poisoned");
*vector_cache = None;
}
let db = Arc::clone(&transaction.db);
let session = transaction.session.clone();
let overlay = &mut transaction.overlay;
let catalog_modified = &mut transaction.catalog_modified;
let mut outcome = Ok(alopex_sql::ExecutionResult::Success);
session
.with_transaction(|owned| {
outcome = (|| {
let mut raw = AnyKVTransaction::Owned(OwnedKVTransactionAdapter::new(owned));
let mode = raw.mode();
let mut borrowed =
TxnBridge::<alopex_core::kv::AnyKV>::wrap_external(&mut raw, mode, overlay);
let mut executor: Executor<_, _> =
Executor::new(db.store.clone(), db.sql_catalog.clone());
let mut last = alopex_sql::ExecutionResult::Success;
for (statement_index, stmt) in statements.iter().enumerate() {
let plan = {
let catalog = db.sql_catalog.read().expect("catalog lock poisoned");
let (_, overlay) = borrowed.split_parts();
plan_stmt(&*catalog, &*overlay, stmt)?
};
{
let catalog = db.sql_catalog.read().expect("catalog lock poisoned");
let (_, overlay) = borrowed.split_parts();
let view = TxnCatalogView::new(&*catalog, &*overlay);
db.record_routing(&view, stmt, statement_index);
}
last = executor
.execute_in_txn(plan, &mut borrowed)
.map_err(|error| Error::Sql(alopex_sql::SqlError::from(error)))?;
}
if statements.iter().any(stmt_changes_catalog) {
*catalog_modified = true;
}
Ok(last)
})();
Ok(())
})
.map_err(Error::Core)?;
outcome
}
fn build_column_info(
projection: &Projection,
schema: &[alopex_sql::catalog::ColumnMetadata],
) -> Result<Vec<ColumnInfo>> {
match projection {
Projection::All(names) => {
let mut cols = Vec::with_capacity(names.len());
for name in names {
let meta = schema.iter().find(|c| &c.name == name).ok_or_else(|| {
Error::Sql(alopex_sql::SqlError::Execution {
message: format!("column not found: {}", name),
code: "ALOPEX-E020",
})
})?;
cols.push(ColumnInfo::new(name.clone(), meta.data_type.clone()));
}
Ok(cols)
}
Projection::Columns(cols) => {
let mut result = Vec::with_capacity(cols.len());
for (i, col) in cols.iter().enumerate() {
let name = col
.alias
.clone()
.or_else(|| {
if let alopex_sql::planner::typed_expr::TypedExprKind::ColumnRef {
column,
..
} = &col.expr.kind
{
Some(column.clone())
} else {
None
}
})
.unwrap_or_else(|| format!("col_{}", i));
result.push(ColumnInfo::new(name, col.expr.resolved_type.clone()));
}
Ok(result)
}
}
}
impl Database {
pub fn begin_read_at_sql(
&self,
point: ReadAtPoint,
) -> Result<alopex_sql::storage::SqlTransaction<'_, alopex_core::kv::AnyKV>> {
let transaction = self.store.begin_read_at(&point).map_err(Error::ReadAt)?;
Ok(TxnBridge::from_read_at(transaction, point))
}
pub fn execute_sql(&self, sql: &str) -> Result<SqlResult> {
Ok(self
.execute_sql_multi(sql)?
.pop()
.unwrap_or(alopex_sql::ExecutionResult::Success))
}
pub fn execute_sql_multi(&self, sql: &str) -> Result<Vec<SqlResult>> {
let stmts = parse_sql(sql)?;
if stmts.is_empty() {
return Ok(Vec::new());
}
if stmts.len() == 1 {
let overlay = CatalogOverlay::new();
let plan = {
let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
plan_stmt(&*catalog, &overlay, &stmts[0])?
};
if matches!(stmts[0].kind, StatementKind::Pragma { .. })
|| alopex_sql::executor::is_store_direct_plan(&plan)
{
let mut executor: Executor<_, _> =
Executor::new(self.store.clone(), self.sql_catalog.clone());
return executor
.execute(plan)
.map(|result| vec![result])
.map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)));
}
}
let requires_write = stmts.iter().any(stmt_requires_write);
let mode = if requires_write {
TxnMode::ReadWrite
} else {
TxnMode::ReadOnly
};
let mut txn = self.store.begin(mode).map_err(Error::Core)?;
let journal = if mode == TxnMode::ReadWrite
&& stmts.iter().any(stmt_changes_user_data)
&& self.store.range_change_journal_capability()
== RangeChangeJournalCapability::Supported
{
let scope = {
let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
local_journal_scope(&*catalog)
};
Some(LocalRangeChangeJournal::capture(&mut txn, scope).map_err(Error::Core)?)
} else {
None
};
let mut overlay = CatalogOverlay::new();
let mut borrowed =
TxnBridge::<alopex_core::kv::AnyKV>::wrap_external(&mut txn, mode, &mut overlay);
let mut executor: Executor<_, _> =
Executor::new(self.store.clone(), self.sql_catalog.clone());
let mut results = Vec::with_capacity(stmts.len());
for (statement_index, stmt) in stmts.iter().enumerate() {
let plan = {
let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
let (_, overlay) = borrowed.split_parts();
plan_stmt(&*catalog, &*overlay, stmt)?
};
{
let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
let (_, overlay) = borrowed.split_parts();
let view = TxnCatalogView::new(&*catalog, &*overlay);
self.record_routing(&view, stmt, statement_index);
}
results.push(
executor
.execute_in_txn(plan, &mut borrowed)
.map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))?,
);
}
drop(borrowed);
if let Some(journal) = journal {
journal.stage(&mut txn).map_err(Error::Core)?;
}
txn.commit_self().map_err(Error::Core)?;
if mode == TxnMode::ReadWrite {
let mut catalog = self.sql_catalog.write().expect("catalog lock poisoned");
catalog.apply_overlay(overlay);
}
if stmts.iter().any(stmt_changes_catalog) {
self.invalidate_table_info_cache();
}
if requires_write {
let mut cache = self.hnsw_cache.write().expect("hnsw cache lock poisoned");
cache.clear();
let mut vector_cache = self
.vector_cache
.write()
.expect("vector cache lock poisoned");
*vector_cache = None;
}
Ok(results)
}
pub fn execute_sql_with_rows<F, R>(&self, sql: &str, f: F) -> Result<StreamingQueryResult<R>>
where
F: FnOnce(StreamingRows<'_>) -> Result<R>,
{
let stmts = parse_sql(sql)?;
if stmts.is_empty() {
return Ok(StreamingQueryResult::Success);
}
if stmts.len() == 1 && matches!(stmts[0].kind, StatementKind::Select(_)) {
let stmt = &stmts[0];
let mode = TxnMode::ReadOnly;
let mut txn = self.store.begin(mode).map_err(Error::Core)?;
let mut overlay = CatalogOverlay::new();
let mut borrowed =
TxnBridge::<alopex_core::kv::AnyKV>::wrap_external(&mut txn, mode, &mut overlay);
let plan = {
let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
let (_, overlay_ref) = borrowed.split_parts();
plan_stmt(&*catalog, overlay_ref, stmt)?
};
if alopex_sql::executor::is_store_direct_plan(&plan) {
let mut executor: Executor<_, _> =
Executor::new(self.store.clone(), self.sql_catalog.clone());
let result = executor
.execute(plan)
.map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))?;
let ExecutionResult::Query(query) = result else {
return Err(Error::Sql(alopex_sql::SqlError::Execution {
message: "store-direct system function did not return rows".into(),
code: "ALOPEX-E022",
}));
};
let column_names: Vec<String> = query
.columns
.iter()
.map(|column| column.name.clone())
.collect();
let schema: Vec<alopex_sql::catalog::ColumnMetadata> = query
.columns
.iter()
.map(|column| {
alopex_sql::catalog::ColumnMetadata::new(
&column.name,
column.data_type.clone(),
)
})
.collect();
let rows: Vec<Row> = query
.rows
.into_iter()
.enumerate()
.map(|(index, values)| Row::new(index as u64, values))
.collect();
let iter = VecIterator::new(rows, schema.clone());
let streaming_rows = StreamingRows {
columns: query.columns,
iter: Box::new(iter),
projection: Projection::All(column_names),
schema,
};
let result = f(streaming_rows)?;
return Ok(StreamingQueryResult::QueryProcessed(result));
}
let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
let (mut sql_txn, overlay_ref) = borrowed.split_parts();
let view = TxnCatalogView::new(&*catalog, overlay_ref);
let (iter, projection, schema) = build_streaming_pipeline(&mut sql_txn, &view, plan)
.map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))?;
let columns = build_column_info(&projection, &schema)?;
let streaming_rows = StreamingRows {
columns,
iter,
projection,
schema,
};
let result = f(streaming_rows)?;
drop(catalog);
drop(borrowed);
txn.commit_self().map_err(Error::Core)?;
return Ok(StreamingQueryResult::QueryProcessed(result));
}
let exec_result = self.execute_sql(sql)?;
match exec_result {
alopex_sql::ExecutionResult::Success => Ok(StreamingQueryResult::Success),
alopex_sql::ExecutionResult::RowsAffected(n) => {
Ok(StreamingQueryResult::RowsAffected(n))
}
alopex_sql::ExecutionResult::Query(_qr) => {
Err(Error::Sql(alopex_sql::SqlError::Execution {
message: "Streaming not available for multi-statement or complex queries"
.into(),
code: "ALOPEX-E021",
}))
}
}
}
pub fn execute_sql_streaming(&self, sql: &str) -> Result<SqlStreamingResult> {
let stmts = parse_sql(sql)?;
if stmts.is_empty() {
return Ok(SqlStreamingResult::Success);
}
if stmts.len() == 1 && matches!(stmts[0].kind, StatementKind::Select(_)) {
let stmt = &stmts[0];
let mode = TxnMode::ReadOnly;
let mut txn = self.store.begin(mode).map_err(Error::Core)?;
let mut overlay = CatalogOverlay::new();
let mut borrowed =
TxnBridge::<alopex_core::kv::AnyKV>::wrap_external(&mut txn, mode, &mut overlay);
let plan = {
let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
let (_, overlay) = borrowed.split_parts();
plan_stmt(&*catalog, &*overlay, stmt)?
};
let (mut sql_txn, _overlay) = borrowed.split_parts();
let catalog = self.sql_catalog.read().expect("catalog lock poisoned");
let view = TxnCatalogView::new(&*catalog, _overlay);
let iter = execute_query_streaming(&mut sql_txn, &view, plan)
.map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))?;
drop(catalog);
drop(borrowed);
txn.commit_self().map_err(Error::Core)?;
return Ok(SqlStreamingResult::Query(iter));
}
let result = self.execute_sql(sql)?;
match result {
alopex_sql::ExecutionResult::Success => Ok(SqlStreamingResult::Success),
alopex_sql::ExecutionResult::RowsAffected(n) => Ok(SqlStreamingResult::RowsAffected(n)),
alopex_sql::ExecutionResult::Query(qr) => {
use alopex_sql::executor::query::iterator::VecIterator;
use alopex_sql::executor::Row;
use alopex_sql::planner::typed_expr::Projection;
let column_names: Vec<String> = qr.columns.iter().map(|c| c.name.clone()).collect();
let schema: Vec<alopex_sql::catalog::ColumnMetadata> = qr
.columns
.iter()
.map(|c| alopex_sql::catalog::ColumnMetadata::new(&c.name, c.data_type.clone()))
.collect();
let rows: Vec<Row> = qr
.rows
.into_iter()
.enumerate()
.map(|(i, values)| Row::new(i as u64, values))
.collect();
let iter = VecIterator::new(rows, schema.clone());
let query_iter =
QueryRowIterator::new(Box::new(iter), Projection::All(column_names), schema);
Ok(SqlStreamingResult::Query(query_iter))
}
}
}
}
impl<'a> Transaction<'a> {
pub fn execute_sql(&mut self, sql: &str) -> Result<SqlResult> {
let stmts = parse_sql(sql)?;
if stmts.is_empty() {
return Ok(alopex_sql::ExecutionResult::Success);
}
if stmts.iter().any(stmt_requires_write) {
let mut cache = self
.db
.hnsw_cache
.write()
.expect("hnsw cache lock poisoned");
cache.clear();
let mut vector_cache = self
.db
.vector_cache
.write()
.expect("vector cache lock poisoned");
*vector_cache = None;
}
let store = self.db.store.clone();
let sql_catalog = self.db.sql_catalog.clone();
let txn = self.inner.as_mut().ok_or(Error::TxnCompleted)?;
let mode = txn.mode();
let mut borrowed =
TxnBridge::<alopex_core::kv::AnyKV>::wrap_external(txn, mode, &mut self.overlay);
let mut executor: Executor<_, _> = Executor::new(store, sql_catalog.clone());
let mut last = alopex_sql::ExecutionResult::Success;
for (statement_index, stmt) in stmts.iter().enumerate() {
let plan = {
let catalog = sql_catalog.read().expect("catalog lock poisoned");
let (_, overlay) = borrowed.split_parts();
plan_stmt(&*catalog, &*overlay, stmt)?
};
{
let catalog = sql_catalog.read().expect("catalog lock poisoned");
let (_, overlay) = borrowed.split_parts();
let view = TxnCatalogView::new(&*catalog, &*overlay);
self.db.record_routing(&view, stmt, statement_index);
}
last = executor
.execute_in_txn(plan, &mut borrowed)
.map_err(|e| Error::Sql(alopex_sql::SqlError::from(e)))?;
}
if stmts.iter().any(stmt_changes_catalog) {
self.catalog_modified = true;
}
Ok(last)
}
}
#[cfg(test)]
mod tests {
use super::*;
use alopex_core::kv::decode_range_change;
use alopex_core::kv::RangeChangePayload;
#[test]
fn auto_commit_stages_sql_row_and_index_changes_before_visibility() {
let db = Database::open_in_memory().unwrap();
db.execute_sql("CREATE TABLE users (id INTEGER PRIMARY KEY, name TEXT);")
.unwrap();
db.execute_sql("CREATE INDEX idx_users_name ON users (name);")
.unwrap();
db.execute_sql("INSERT INTO users (id, name) VALUES (1, 'alice');")
.unwrap();
let mut reader = db.store.begin(TxnMode::ReadOnly).unwrap();
let records = reader
.scan_prefix(b"\x00alopex/range-change/")
.unwrap()
.filter_map(|(_, value)| decode_range_change(&value).ok())
.collect::<Vec<_>>();
assert_eq!(records.len(), 1);
assert!(records[0]
.payload
.iter()
.any(|payload| matches!(payload, RangeChangePayload::UpsertRow { .. })));
assert!(records[0]
.payload
.iter()
.any(|payload| matches!(payload, RangeChangePayload::UpsertIndex { .. })));
}
#[test]
fn explicit_sql_transaction_stages_journal_before_commit() {
let db = Database::open_in_memory().unwrap();
db.execute_sql("CREATE TABLE users (id INTEGER PRIMARY KEY, name TEXT);")
.unwrap();
let mut transaction = db.begin(TxnMode::ReadWrite).unwrap();
transaction
.execute_sql("INSERT INTO users (id, name) VALUES (1, 'alice');")
.unwrap();
transaction.commit().unwrap();
let records = db
.snapshot()
.into_iter()
.filter_map(|(_, value)| decode_range_change(&value).ok())
.collect::<Vec<_>>();
assert_eq!(records.len(), 1);
assert!(records[0]
.payload
.iter()
.any(|payload| matches!(payload, RangeChangePayload::UpsertRow { .. })));
}
#[test]
fn embedded_read_at_returns_retention_error_before_a_sql_session_exists() {
let db = Database::open_in_memory().unwrap();
let point = ReadAtPoint::new(1, 2, 3, 4);
assert!(matches!(
db.begin_read_at_sql(point),
Err(Error::ReadAt(alopex_core::ReadAtError::Unavailable { .. }))
));
}
}