use super::Connection;
use super::plan_cache::{CachedPlan, normalize_query};
use crate::database::Database;
use crate::prepared_statement::PreparedStatement;
use crate::query_result::QueryResult;
use akar_binder::Binder;
use akar_binder::bound_statement::BoundStatement;
use akar_common::error::ProcessorError;
use akar_common::types::Value;
use akar_optimizer::Optimizer;
use akar_parser::parse;
use akar_planner::QueryPlanner;
use akar_planner::logical_operator::LogicalOperator;
use akar_processor::QueryProcessor;
use akar_processor::processor::{SchemaDdlFn, SchemaDdlOp, SequenceFn, StandaloneCallHandler, SubqueryFn};
use std::collections::HashMap;
use std::sync::Arc;
impl Connection {
pub fn query(&self, query_str: &str) -> Result<QueryResult, String> {
let trimmed = query_str.trim();
if trimmed.is_empty() {
return Ok(QueryResult::new(Vec::new()));
}
if let Some(value) = trimmed
.strip_prefix("SET")
.and_then(|s| s.trim().strip_prefix("spill_threshold"))
.and_then(|s| s.trim().strip_prefix("="))
.map(|s| s.trim())
{
let bytes: u64 = value.parse().map_err(|_| {
format!("Invalid spill_threshold value '{value}'. Expected a non-negative integer (bytes).")
})?;
self.database.set_spill_threshold(bytes);
return Ok(QueryResult::success_message(format!(
"spill_threshold set to {bytes} bytes"
)));
}
if let Some(value) = trimmed
.strip_prefix("SET")
.and_then(|s| s.trim().strip_prefix("concurrent_writes"))
.and_then(|s| s.trim().strip_prefix("="))
.map(|s| s.trim())
{
let enabled = match value.to_lowercase().as_str() {
"true" | "1" | "yes" => true,
"false" | "0" | "no" => false,
_ => return Err("Invalid value for concurrent_writes. Use true or false.".into()),
};
self.database.transaction_manager.set_concurrent_writes(enabled);
return Ok(QueryResult::success_message(format!(
"concurrent_writes set to {enabled}"
)));
}
let normalized = normalize_query(trimmed);
let catalog_version = self
.database
.catalog
.lock()
.map_err(|e| format!("Catalog lock error: {e}"))?
.version();
{
let mut cache = self.plan_cache.lock().map_err(|e| format!("Lock poisoned: {e}"))?;
if let Some(cached) = cache.get(&normalized).filter(|c| c.catalog_version == catalog_version) {
let bound = cached.bound.clone();
let plan = cached.plan.clone();
drop(cache);
return self.execute_with_plan(&bound, Some(&plan));
}
}
let statement = parse(trimmed).map_err(|e| format!("Parse error: {e}"))?;
let binder = Binder::new(self.database.catalog.clone());
let bound = binder.bind(statement).map_err(|e| format!("Bind error: {e}"))?;
let plan_opt: Option<Arc<Vec<LogicalOperator>>> = if is_plan_cachable(&bound) {
let plan = Arc::new(self.build_optimized_plan(&bound)?);
let mut cache = self.plan_cache.lock().map_err(|e| format!("Lock poisoned: {e}"))?;
cache.insert(
normalized,
CachedPlan {
bound: Arc::new(bound.clone()),
plan: Arc::clone(&plan),
catalog_version,
},
);
Some(plan)
} else {
None
};
self.execute_with_plan(&bound, plan_opt.as_ref())
}
fn execute_with_plan(
&self,
bound: &BoundStatement,
plan: Option<&Arc<Vec<LogicalOperator>>>,
) -> Result<QueryResult, String> {
if self.database.config.read_only && Connection::is_write_statement(bound) {
return Err("Database is in read-only mode; write statements are not allowed".into());
}
let is_write = Connection::is_write_statement(bound);
let mut txn_opt: Option<akar_transaction::Transaction> =
if is_write { Some(self.begin_write_txn()?) } else { None };
let query_result = self.execute_query_inner(bound, txn_opt.as_mut(), plan);
match (is_write, &query_result) {
(true, Ok(_)) => {
if let Some(ref mut txn) = txn_opt {
self.commit_write_txn(txn)?;
}
}
(true, Err(e)) => {
if let Some(ref mut txn) = txn_opt {
match self.rollback_write_txn(txn) {
Ok(_records) => {
tracing::warn!("Transaction rolled back due to error: {e}");
}
Err(rollback_err) => {
tracing::error!("Transaction rollback ALSO failed: {rollback_err} (original error: {e})");
}
}
}
}
_ => {}
}
if query_result.is_ok() && Connection::is_write_statement(bound) {
let written = Connection::extract_write_tables(bound);
if !written.is_empty() {
self.database.refresh_vector_indexes(&written);
}
}
query_result
}
fn build_optimized_plan(&self, bound: &BoundStatement) -> Result<Vec<LogicalOperator>, String> {
let planner = QueryPlanner::new();
let logical_plan = planner.plan(bound.clone()).map_err(|e| format!("Plan error: {e}"))?;
let optimizer = Optimizer::with_stats(self.database.stats_store.clone());
let optimized = optimizer.optimize(logical_plan);
Ok(optimized)
}
pub(crate) fn execute_query_inner(
&self,
bound: &BoundStatement,
mut txn_opt: Option<&mut akar_transaction::Transaction>,
cached_plan: Option<&Arc<Vec<LogicalOperator>>>,
) -> Result<QueryResult, String> {
if let Some(result) = self.handle_ddl(bound, txn_opt.as_deref_mut())? {
self.database.persist_catalog()?;
self.maybe_auto_checkpoint()?;
return Ok(result);
}
if let Some(ref txn) = txn_opt {
if !self.database.transaction_manager.allow_concurrent_writes() {
let write_tables = Connection::extract_write_tables(bound);
for tid in write_tables {
self.database.transaction_manager.lock_table(txn.transaction_id, tid)?;
}
}
}
let optimized_plan: Arc<Vec<LogicalOperator>> = match cached_plan {
Some(plan) => plan.clone(),
None => Arc::new(self.build_optimized_plan(bound)?),
};
if optimized_plan.is_empty() {
return Ok(QueryResult::success_message("Query executed (no result)".into()));
}
let (snapshot_ts, commit_history) = if let Some(ref txn) = txn_opt {
(
txn.snapshot_ts,
self.database.transaction_manager.commit_history_snapshot(),
)
} else {
let ts = self.database.transaction_manager.current_commit_ts();
let history = self.database.transaction_manager.commit_history_snapshot();
(Some(ts), history)
};
let processor = self
.create_processor()
.with_snapshot(snapshot_ts, commit_history)
.with_txn_id(txn_opt.as_ref().map(|t| t.transaction_id));
let chunks = processor
.execute(&optimized_plan)
.map_err(|e| format!("Execute error: {e}"))?;
if let Some(ref txn) = txn_opt {
let written_rows = processor.take_written_rows();
let tm = &self.database.transaction_manager;
for (table_id, row_id) in written_rows {
tm.record_write(txn.transaction_id, table_id, row_id);
}
}
if let Some(ref mut txn) = txn_opt {
let undo = processor.take_undo_records();
txn.undo_records.extend(undo);
}
self.maybe_auto_checkpoint()?;
Ok(QueryResult::new(chunks))
}
pub fn prepare(&self, query_str: &str) -> Result<PreparedStatement, String> {
let trimmed = query_str.trim();
{
let cache = self.statement_cache.lock().map_err(|e| format!("Lock poisoned: {e}"))?;
if let Some(cached) = cache.get(trimmed) {
return Ok(cached.clone());
}
}
let statement = parse(trimmed).map_err(|e| format!("Parse error: {e}"))?;
let binder = Binder::new(self.database.catalog.clone());
let bound = binder.bind(statement).map_err(|e| format!("Bind error: {e}"))?;
let prepared = PreparedStatement::new(trimmed.to_string(), bound);
{
let mut cache = self.statement_cache.lock().map_err(|e| format!("Lock poisoned: {e}"))?;
cache.insert(trimmed.to_string(), prepared.clone());
}
Ok(prepared)
}
pub fn execute(&self, prepared: &PreparedStatement, params: Vec<(&str, Value)>) -> Result<QueryResult, String> {
let mut param_map = HashMap::new();
let num_expected = prepared.parameters.len();
for (name, value) in ¶ms {
param_map.insert(name.to_string(), value.clone());
}
for p in &prepared.parameters {
if !param_map.contains_key(p) {
return Err(format!("Missing parameter: ${}", p));
}
}
if params.len() > num_expected {
return Err(format!("Expected {} parameter(s), got {}", num_expected, params.len()));
}
if self.database.config.read_only && Connection::is_write_statement(&prepared.bound_statement) {
return Err("Database is in read-only mode; write statements are not allowed".into());
}
let substituted =
crate::connection::substitute::substitute_params_in_statement(&prepared.bound_statement, ¶m_map)?;
if let Some(result) = self.handle_ddl(&substituted, None)? {
self.database.persist_catalog()?;
self.maybe_auto_checkpoint()?;
return Ok(result);
}
let planner = QueryPlanner::new();
let logical_plan = planner.plan(substituted).map_err(|e| format!("Plan error: {e}"))?;
if logical_plan.is_empty() {
return Ok(QueryResult::success_message("Query executed (no result)".into()));
}
let optimizer = Optimizer::with_stats(self.database.stats_store.clone());
let optimized_plan = optimizer.optimize(logical_plan);
let is_write = Connection::is_write_statement(&prepared.bound_statement);
let mut txn_opt: Option<akar_transaction::Transaction> =
if is_write { Some(self.begin_write_txn()?) } else { None };
let (snapshot_ts, history) = if let Some(ref txn) = txn_opt {
(
txn.snapshot_ts,
self.database.transaction_manager.commit_history_snapshot(),
)
} else {
let ts = self.database.transaction_manager.current_commit_ts();
(Some(ts), self.database.transaction_manager.commit_history_snapshot())
};
let processor = self
.create_processor()
.with_snapshot(snapshot_ts, history)
.with_txn_id(txn_opt.as_ref().map(|t| t.transaction_id));
let chunks = match processor.execute(&optimized_plan) {
Ok(c) => c,
Err(e) => {
if is_write {
if let Some(ref mut txn) = txn_opt {
match self.rollback_write_txn(txn) {
Ok(_) => tracing::warn!("Prepared write rolled back due to error: {e}"),
Err(rollback_err) => {
tracing::error!("Prepared rollback ALSO failed: {rollback_err} (original: {e})");
}
}
}
}
return Err(format!("Execute error: {e}"));
}
};
if let Some(ref txn) = txn_opt {
let written_rows = processor.take_written_rows();
let tm = &self.database.transaction_manager;
for (table_id, row_id) in written_rows {
tm.record_write(txn.transaction_id, table_id, row_id);
}
}
if let Some(ref mut txn) = txn_opt {
let undo = processor.take_undo_records();
txn.undo_records.extend(undo);
}
if is_write {
if let Some(ref mut txn) = txn_opt {
self.commit_write_txn(txn)?;
}
}
self.maybe_auto_checkpoint()?;
Ok(QueryResult::new(chunks))
}
pub(crate) fn maybe_auto_checkpoint(&self) -> Result<(), String> {
let threshold = self.database.config.checkpoint_threshold;
if threshold == 0 {
return Ok(()); }
let should_checkpoint = if threshold < 0 {
true
} else {
self.database.storage_manager.wal_size() > threshold as usize
};
if should_checkpoint {
self.database.transaction_manager.schedule_auto_checkpoint();
tracing::debug!("Auto-checkpoint signaled to background worker");
}
Ok(())
}
pub(crate) fn do_sync_checkpoint(&self) -> Result<(), String> {
let tm = &self.database.transaction_manager;
let drain_fn = |timeout: std::time::Duration| -> bool { tm.stop_new_txns_and_wait_until_all_leave(timeout) };
self.database
.storage_manager
.checkpoint_with_drain(Some(&drain_fn))
.map_err(|e| format!("Checkpoint failed: {e}"))?;
tracing::debug!("Sync checkpoint completed");
Ok(())
}
pub(crate) fn create_processor(&self) -> QueryProcessor {
let handlers = self
.processor_handlers
.get_or_init(|| Arc::new(build_processor_handlers(&self.database)));
QueryProcessor::with_catalog(
self.database.function_registry.clone(),
self.database.table_catalog(),
self.database.vfs.clone(),
)
.with_sequence_fn(handlers.sequence_fn.clone())
.with_subquery_fn(handlers.subquery_fn.clone())
.with_schema_ddl_fn(handlers.schema_ddl_fn.clone())
.with_standalone_call_handler(handlers.standalone_call_handler.clone())
}
}
pub(crate) struct ProcessorHandlers {
pub sequence_fn: SequenceFn,
pub schema_ddl_fn: SchemaDdlFn,
pub subquery_fn: SubqueryFn,
pub standalone_call_handler: Arc<dyn StandaloneCallHandler>,
}
fn build_processor_handlers(db: &Arc<Database>) -> ProcessorHandlers {
let seq_fn = super::utils::make_sequence_callback(db.catalog.clone());
let db_sddl = db.clone();
let schema_ddl_fn: SchemaDdlFn = Arc::new(move |op: SchemaDdlOp| -> Result<String, ProcessorError> {
match op {
SchemaDdlOp::CreateSequence {
name,
if_not_exists,
start_value,
increment,
min_value,
max_value,
cycle,
} => {
let mut catalog = db_sddl.catalog.lock().map_err(|e| format!("Catalog lock: {e}"))?;
match catalog.create_sequence(name.clone(), start_value, increment, min_value, max_value, cycle) {
akar_catalog::CatalogResult::Created { .. } => Ok(format!("Sequence '{}' created", name)),
akar_catalog::CatalogResult::AlreadyExists => {
if if_not_exists {
Ok(format!("Sequence '{}' already exists", name))
} else {
Err(ProcessorError::Execution(format!("Sequence '{}' already exists", name)))
}
}
other => Err(ProcessorError::Execution(format!(
"Failed to create sequence: {:?}",
other
))),
}
}
SchemaDdlOp::DropSequence { name, if_exists } => {
let mut catalog = db_sddl.catalog.lock().map_err(|e| format!("Catalog lock: {e}"))?;
match catalog.drop_sequence(&name) {
akar_catalog::CatalogResult::Dropped { .. } => Ok(format!("Sequence '{}' dropped", name)),
akar_catalog::CatalogResult::NotFound => {
if if_exists {
Ok(format!("Sequence '{}' not found", name))
} else {
Err(ProcessorError::Execution(format!("Sequence '{}' not found", name)))
}
}
other => Err(ProcessorError::Execution(format!(
"Failed to drop sequence: {:?}",
other
))),
}
}
SchemaDdlOp::ExportDatabase {
file_path,
file_type,
schema_only,
} => {
let conn = super::Connection::new(&db_sddl);
let bound = akar_binder::bound_statement::BoundExportDatabase {
file_path,
file_type,
schema_only,
options: Default::default(),
};
let result = conn
.execute_export_database(&bound)
.map_err(ProcessorError::Execution)?;
let msg = result
.and_then(|r| r.message)
.unwrap_or_else(|| format!("Database exported to '{}'", bound.file_path));
Ok(msg)
}
SchemaDdlOp::ImportDatabase {
file_path,
query,
index_query,
} => {
let conn = super::Connection::new(&db_sddl);
let mut executed = 0usize;
let mut skipped = 0usize;
for stmt in super::copy::split_cypher_statements(&query)
.into_iter()
.chain(super::copy::split_cypher_statements(&index_query))
{
match conn.query(&stmt) {
Ok(_) => executed += 1,
Err(e) => {
tracing::warn!("Import statement skipped (may be duplicate): {e}");
skipped += 1;
}
}
}
Ok(format!(
"Imported {executed} statement(s) from '{file_path}' ({skipped} skipped)"
))
}
}
});
let db_qf = db.clone();
let query_fn: crate::connection::standalone_call::QueryFn = Arc::new({
let schema_ddl_qf = schema_ddl_fn.clone();
move |query_str: &str| -> Result<crate::query_result::QueryResult, String> {
let stmt = akar_parser::parse(query_str).map_err(|e| format!("Parse error: {e}"))?;
let binder = Binder::new(db_qf.catalog.clone());
let bound = binder.bind(stmt).map_err(|e| format!("Bind error: {e}"))?;
let planner = QueryPlanner::new();
let logical_plan = planner.plan(bound).map_err(|e| format!("Plan error: {e}"))?;
let optimizer = Optimizer::with_stats(db_qf.stats_store.clone());
let optimized_plan = optimizer.optimize(logical_plan);
let processor = QueryProcessor::with_catalog(
db_qf.function_registry.clone(),
db_qf.table_catalog(),
db_qf.vfs.clone(),
)
.with_schema_ddl_fn(schema_ddl_qf.clone())
.with_standalone_call_handler(Arc::new(
crate::connection::standalone_call::DbStandaloneCallHandler::new(db_qf.clone()),
))
.with_snapshot(
Some(db_qf.transaction_manager.current_commit_ts()),
db_qf.transaction_manager.commit_history_snapshot(),
);
let chunks = processor
.execute(&optimized_plan)
.map_err(|e| format!("Execute error: {e}"))?;
let num_rows: usize = chunks.iter().map(|c| c.size).sum();
let num_columns = chunks.first().map(|c| c.num_fields()).unwrap_or(0);
Ok(crate::query_result::QueryResult {
chunks,
num_rows,
num_columns,
success: true,
error_message: None,
message: None,
summary: None,
})
}
});
let db_sq = db.clone();
let subquery_fn: SubqueryFn = Arc::new({
let schema_ddl_sq = schema_ddl_fn.clone();
move |query: &akar_parser::ast::Query| -> Result<Vec<akar_common::vector::DataChunk>, ProcessorError> {
let stmt = akar_parser::ast::Statement::Query(query.clone());
let binder = Binder::new(db_sq.catalog.clone());
let bound = binder.bind(stmt).map_err(|e| format!("Bind error: {e}"))?;
let planner = QueryPlanner::new();
let logical_plan = planner.plan(bound).map_err(|e| format!("Plan error: {e}"))?;
let optimizer = Optimizer::with_stats(db_sq.stats_store.clone());
let optimized_plan = optimizer.optimize(logical_plan);
let catalog_inner = db_sq.catalog.clone();
let seq_fn_inner = super::utils::make_sequence_callback(catalog_inner);
let processor = QueryProcessor::with_catalog(
db_sq.function_registry.clone(),
db_sq.table_catalog(),
db_sq.vfs.clone(),
)
.with_sequence_fn(seq_fn_inner)
.with_schema_ddl_fn(schema_ddl_sq.clone())
.with_standalone_call_handler(Arc::new(
crate::connection::standalone_call::DbStandaloneCallHandler::new(db_sq.clone()),
))
.with_snapshot(
Some(db_sq.transaction_manager.current_commit_ts()),
db_sq.transaction_manager.commit_history_snapshot(),
);
processor
.execute(&optimized_plan)
.map_err(|e| ProcessorError::Execution(format!("Execute error: {e}")))
}
});
let standalone_call_handler: Arc<dyn StandaloneCallHandler> = Arc::new(
crate::connection::standalone_call::DbStandaloneCallHandler::with_query_executor(
db.clone(),
Some(query_fn.clone()),
),
);
ProcessorHandlers {
sequence_fn: seq_fn,
schema_ddl_fn,
subquery_fn,
standalone_call_handler,
}
}
fn is_plan_cachable(bound: &BoundStatement) -> bool {
match bound {
BoundStatement::BoundQuery(q) => {
!(q.clauses.len() == 1
&& matches!(
q.clauses.first(),
Some(akar_binder::bound_statement::BoundClause::BoundForeach(_))
))
}
_ => false,
}
}