use std::any::Any;
use std::collections::BTreeMap;
use std::sync::{Arc, Mutex};
use anyhow::{anyhow, Result};
use lora_ast::{Direction, Document};
use lora_executor::{lora_value_to_property, ExecuteOptions, LoraValue, QueryResult};
use lora_parser::parse_query;
use lora_store::{GraphStorage, GraphStorageMut, InMemoryGraph, Properties};
use lora_wal::WalRecorder;
mod builder;
mod compile;
mod execute;
mod explain;
mod graph_api;
mod occ;
mod procedures;
mod profile;
mod pull_mode;
mod replay;
mod row_projection;
mod schema;
mod show_pipeline;
mod stream;
mod write_guard;
use crate::error::LoraError;
use crate::explain::{QueryPlan, QueryProfile};
use crate::live_store::LiveStore;
use crate::plan_cache::PlanCache;
use crate::snapshot::ManagedSnapshotStore;
use crate::wal::archive::WalArchive;
pub trait QueryRunner: Send + Sync + 'static {
fn execute(
&self,
query: &str,
options: Option<ExecuteOptions>,
) -> Result<QueryResult, LoraError>;
fn explain(
&self,
query: &str,
params: Option<BTreeMap<String, LoraValue>>,
) -> Result<QueryPlan, LoraError>;
fn profile(
&self,
query: &str,
params: Option<BTreeMap<String, LoraValue>>,
) -> Result<QueryProfile, LoraError>;
}
pub struct Database<S> {
pub(crate) store: Arc<LiveStore<S>>,
pub(crate) writer: Arc<Mutex<()>>,
#[allow(dead_code)]
pub(crate) lock_table: Arc<lora_store::LockTable>,
pub(crate) wal: Option<Arc<WalRecorder>>,
pub(crate) snapshots: Option<Arc<ManagedSnapshotStore>>,
pub(crate) named_archive: Option<Arc<WalArchive>>,
pub(crate) plan_cache: Arc<PlanCache>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum GraphDirection {
Outgoing,
Incoming,
Both,
}
impl GraphDirection {
pub(crate) fn as_store_direction(self) -> Direction {
match self {
Self::Outgoing => Direction::Right,
Self::Incoming => Direction::Left,
Self::Both => Direction::Undirected,
}
}
}
pub(crate) fn values_to_properties(values: BTreeMap<String, LoraValue>) -> Result<Properties> {
values
.into_iter()
.map(|(key, value)| {
let value = lora_value_to_property(value).map_err(|e| anyhow!(e))?;
Ok((key, value))
})
.collect()
}
pub(crate) const QUERY_FAILURE_POISON: &str =
"query mutated the live graph before failing; restart from snapshot + WAL required";
impl Database<InMemoryGraph> {
pub fn sync(&self) -> Result<(), LoraError> {
if let Some(wal) = &self.wal {
if let Some(archive) = &self.named_archive {
let _commit_lock = self
.writer
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
wal.force_fsync_wal_only()?;
let snapshot_lsn = wal.wal().durable_lsn();
let graph = self.store.load_full();
let payload = graph.snapshot_payload();
let mut bytes = Vec::new();
let options = lora_snapshot::SnapshotOptions::default();
crate::snapshot::encode_snapshot_to(
&mut bytes,
&payload,
Some(snapshot_lsn.raw()),
&options,
)?;
archive.persist_snapshot_bytes(bytes)?;
} else {
wal.force_fsync()?;
}
}
Ok(())
}
}
impl<S> Database<S>
where
S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
{
pub fn wal(&self) -> Option<&Arc<WalRecorder>> {
self.wal.as_ref()
}
pub fn snapshot(&self) -> Arc<S> {
self.store.load_full()
}
pub fn parse(&self, query: &str) -> Result<Document, LoraError> {
Ok(parse_query(query)?)
}
pub(crate) fn read_store(&self) -> Arc<S> {
self.store.load_full()
}
}
impl<S> Database<S>
where
S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
{
pub fn try_clear(&self) -> Result<(), LoraError> {
self.with_logged_store_mut(|store| {
store.clear();
Ok(())
})
.map_err(LoraError::from_anyhow)
}
pub fn clear(&self) {
let _ = self.try_clear();
}
pub fn node_count(&self) -> usize {
let snapshot = self.read_store();
snapshot.node_count()
}
pub fn relationship_count(&self) -> usize {
let snapshot = self.read_store();
snapshot.relationship_count()
}
pub fn with_store<R>(&self, f: impl FnOnce(&S) -> R) -> R {
let snapshot = self.read_store();
f(&*snapshot)
}
pub fn with_store_mut<R>(&self, f: impl FnOnce(&mut S) -> R) -> R {
let _lock = self
.writer
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
let mut handle = self.store.write();
f(handle.as_mut())
}
}
impl<S> QueryRunner for Database<S>
where
S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
{
fn execute(
&self,
query: &str,
options: Option<ExecuteOptions>,
) -> Result<QueryResult, LoraError> {
Database::execute(self, query, options)
}
fn explain(
&self,
query: &str,
params: Option<BTreeMap<String, LoraValue>>,
) -> Result<QueryPlan, LoraError> {
Database::explain(self, query, params)
}
fn profile(
&self,
query: &str,
params: Option<BTreeMap<String, LoraValue>>,
) -> Result<QueryProfile, LoraError> {
Database::profile(self, query, params)
}
}