use std::any::Any;
use std::collections::BTreeMap;
use std::time::{Duration, Instant};
use anyhow::Result;
use lora_executor::{
classify_stream, collect_compiled, project_rows, ExecuteOptions, LoraValue,
MutableExecutionContext, MutableExecutor, QueryResult, Row, StreamShape,
};
use lora_store::{GraphStorage, GraphStorageMut, InMemoryGraph};
use crate::database::{Database, QUERY_FAILURE_POISON};
use crate::error::LoraError;
use crate::wal::write_scope::{ensure_wal_query_can_start, WalAbortPolicy};
use super::pull_mode::should_collect_read_via_pull;
impl<S> Database<S>
where
S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
{
pub fn execute(
&self,
query: &str,
options: Option<ExecuteOptions>,
) -> Result<QueryResult, LoraError> {
self.execute_with_params(query, options, BTreeMap::new())
}
pub fn execute_with_timeout(
&self,
query: &str,
options: Option<ExecuteOptions>,
timeout: Duration,
) -> Result<QueryResult, LoraError> {
let deadline = Instant::now()
.checked_add(timeout)
.unwrap_or_else(Instant::now);
let rows =
self.execute_rows_with_params_deadline(query, BTreeMap::new(), Some(deadline))?;
Ok(project_rows(rows, options.unwrap_or_default()))
}
pub fn execute_with_params(
&self,
query: &str,
options: Option<ExecuteOptions>,
params: BTreeMap<String, LoraValue>,
) -> Result<QueryResult, LoraError> {
let rows = self.execute_rows_with_params_deadline(query, params, None)?;
Ok(project_rows(rows, options.unwrap_or_default()))
}
pub fn execute_with_params_timeout(
&self,
query: &str,
options: Option<ExecuteOptions>,
params: BTreeMap<String, LoraValue>,
timeout: Duration,
) -> Result<QueryResult, LoraError> {
let deadline = Instant::now()
.checked_add(timeout)
.unwrap_or_else(Instant::now);
let rows = self.execute_rows_with_params_deadline(query, params, Some(deadline))?;
Ok(project_rows(rows, options.unwrap_or_default()))
}
pub fn execute_rows(&self, query: &str) -> Result<Vec<Row>, LoraError> {
self.execute_rows_with_params(query, BTreeMap::new())
}
pub fn execute_rows_with_params(
&self,
query: &str,
params: BTreeMap<String, LoraValue>,
) -> Result<Vec<Row>, LoraError> {
Ok(self.execute_rows_with_params_deadline(query, params, None)?)
}
pub(crate) fn execute_rows_with_params_deadline(
&self,
query: &str,
params: BTreeMap<String, LoraValue>,
deadline: Option<Instant>,
) -> Result<Vec<Row>> {
let store = self.read_store_deadline(deadline)?;
let compiled = self.compile_query_cached(query, &*store)?;
let shape = classify_stream(&compiled);
if matches!(shape, StreamShape::ReadOnly) {
if let Some(rec) = &self.wal {
ensure_wal_query_can_start(rec)?;
}
if deadline.is_none() && should_collect_read_via_pull(&compiled) {
return collect_compiled(&*store, params, &compiled).map_err(anyhow::Error::from);
}
let executor = lora_executor::Executor::with_deadline(
lora_executor::ExecutionContext {
storage: &*store,
params,
},
deadline,
);
return executor
.execute_compiled_rows(&compiled)
.map_err(anyhow::Error::from);
}
drop(store);
debug_assert!(shape.is_mutating());
if std::any::TypeId::of::<S>() == std::any::TypeId::of::<InMemoryGraph>() {
return self.execute_mutating_optimistic(params, deadline, &compiled);
}
let store = self.write_store_deadline(deadline)?;
self.with_logged_write_guard(
store,
WalAbortPolicy::PoisonIfMutated(QUERY_FAILURE_POISON),
|store| {
let mut executor = MutableExecutor::with_deadline(
MutableExecutionContext {
storage: store,
params,
},
deadline,
);
Ok(executor.execute_compiled_rows(&compiled)?)
},
)
}
}