use crate::store::Store;
use crate::AsyncStore;
use oxilite_core::job::Job;
use oxilite_core::{AsyncBackend, Step, SyncBackend};
use oxilite_cypher::{
prepare_for, schema_command, CypherError, CypherOptions, CypherResult, CypherStep, Params,
Schema, SchemaCommand, SqlCypherJob, StepInput, Value,
};
use std::borrow::Cow;
impl<B: SyncBackend + Send + Sync + 'static> Store<B> {
pub fn cypher(&self, query: &str) -> Result<CypherResult, CypherError> {
self.cypher_with(query, &Params::new(), &CypherOptions::default())
}
pub fn cypher_with(
&self,
query: &str,
params: &Params,
options: &CypherOptions,
) -> Result<CypherResult, CypherError> {
if let Some(command) = schema_command(query)? {
return self.cypher_schema_command(command, options);
}
let options = &*self.cypher_functions(options);
let resolved;
let options = match &options.query.as_of {
Some(v) if options.query.as_of_tick.is_none() => {
let mut o = options.clone();
o.query.as_of_tick = Some(self.resolve_version(v)?);
resolved = o;
&resolved
}
_ => options,
};
let mut job = prepare_for(query, params, options, self.caps())?;
if job.writes() && options.query.as_of.is_some() {
return Err(versioned_write());
}
let backend = self.versioned();
let tx = job.writes() && self.caps().interactive_transactions;
if tx {
backend.begin()?;
}
let mut run = || -> Result<CypherResult, CypherError> {
let mut input = None;
loop {
input = Some(match job.step(input)? {
CypherStep::Query(q, o) => StepInput::Output(self.evaluate(&q, &o)?),
CypherStep::Sql(r) => StepInput::Response(backend.execute(&r)?),
CypherStep::Write(r) => StepInput::Response(backend.execute(&r)?),
CypherStep::Done(r) => return Ok(r),
});
}
};
match run() {
Ok(r) => {
if tx {
backend.commit()?;
}
if r.schema_changed {
self.reload_stats()?;
}
Ok(r)
}
Err(e) => {
if tx {
let _ = backend.rollback();
}
Err(e)
}
}
}
pub fn cypher_schema(&self) -> Result<std::sync::Arc<Schema>, CypherError> {
let index = self.shape_index().map_err(CypherError::Store)?;
Ok(std::sync::Arc::new(Schema::from_index(index)))
}
fn cypher_functions<'o>(&self, options: &'o CypherOptions) -> Cow<'o, CypherOptions> {
let functions = self.host_functions();
if options.functions.registry().is_some() || functions.is_empty() {
Cow::Borrowed(options)
} else {
let mut o = options.clone();
o.functions = functions;
Cow::Owned(o)
}
}
fn cypher_schema_command(
&self,
command: SchemaCommand,
options: &CypherOptions,
) -> Result<CypherResult, CypherError> {
use oxilite_core::vector::{ElementType, Metric, VectorIndex};
let vocab = &options.vocabulary;
let empty = |columns: Vec<String>, rows: Vec<Vec<Value>>| CypherResult {
columns,
rows,
stats: Default::default(),
schema_changed: false,
};
match command {
SchemaCommand::CreateVectorIndex {
name,
if_not_exists,
label,
property,
dimensions,
similarity,
element_type,
} => {
let exists = self
.stats()
.vector_indexes
.iter()
.any(|d| d.name.eq_ignore_ascii_case(&name));
if exists && if_not_exists {
return Ok(empty(Vec::new(), Vec::new()));
}
let dimensions = dimensions.ok_or_else(|| {
CypherError::Semantic(format!(
"vector index {name} needs OPTIONS {{indexConfig: {{`vector.dimensions`: N}}}}"
))
})?;
let mut index = VectorIndex::new(name, vocab.iri(&property), dimensions);
if let Some(s) = similarity {
index.metric = Metric::parse(&s).ok_or_else(|| {
CypherError::Semantic(format!(
"unknown vector.similarity_function {s:?} (cosine, euclidean, dot, jaccard)"
))
})?;
if index.metric == Metric::Jaccard {
index.element_type = ElementType::SparseFloat32;
}
}
if let Some(t) = element_type {
index.element_type = ElementType::parse(&t).ok_or_else(|| {
CypherError::Semantic(format!(
"unknown vector.element_type {t:?} (float32, float64, int8, bit1, sparse)"
))
})?;
}
if let Some(l) = label {
index.class = Some(vocab.iri(&l));
}
self.create_vector_index(&index)?;
Ok(empty(Vec::new(), Vec::new()))
}
SchemaCommand::DropIndex { name, if_exists } => {
if !self.drop_vector_index(&name)? && !if_exists {
return Err(CypherError::Semantic(format!(
"there is no index named {name}"
)));
}
Ok(empty(Vec::new(), Vec::new()))
}
SchemaCommand::ShowVectorIndexes => {
let columns = [
"name",
"label",
"property",
"dimensions",
"similarityFunction",
"elementType",
"rows",
"state",
]
.map(String::from)
.to_vec();
let rows = self
.vector_indexes()?
.into_iter()
.map(|i| {
vec![
Value::String(i.index.name.clone()),
i.index
.class
.as_ref()
.map_or(Value::Null, |c| Value::String(vocab.name(c.as_str()))),
Value::String(vocab.name(i.index.property.as_str())),
Value::Int(i64::from(i.index.dimensions)),
Value::String(i.index.metric.name().into()),
Value::String(i.index.element_type.name().into()),
Value::Int(i.rows as i64),
Value::String(if i.built { "ONLINE" } else { "NOT BUILT" }.into()),
]
})
.collect();
Ok(empty(columns, rows))
}
}
}
pub fn explain_cypher(
&self,
query: &str,
params: &Params,
options: &CypherOptions,
) -> Result<String, CypherError> {
let options = &*self.cypher_functions(options);
let job = prepare_for(query, params, options, self.caps())?;
let mut out = job.explain();
for (q, o) in job.queries() {
out.push_str(&self.explain_opt(q, &o)?);
out.push('\n');
}
Ok(out)
}
}
impl<B: AsyncBackend> AsyncStore<B> {
pub async fn cypher(&self, query: &str) -> Result<CypherResult, CypherError> {
self.cypher_with(query, &Params::new(), &CypherOptions::default())
.await
}
pub async fn cypher_schema(&self) -> Result<std::sync::Arc<Schema>, CypherError> {
let index = self.shape_index().await.map_err(CypherError::Store)?;
Ok(std::sync::Arc::new(Schema::from_index(index)))
}
pub async fn cypher_with(
&self,
query: &str,
params: &Params,
options: &CypherOptions,
) -> Result<CypherResult, CypherError> {
if schema_command(query)?.is_some() {
return Err(CypherError::Unsupported(
"vector index commands run on the blocking store (Store::cypher)".into(),
));
}
let resolved;
let options = match &options.query.as_of {
Some(v) if options.query.as_of_tick.is_none() => {
let mut o = options.clone();
o.query.as_of_tick = Some(self.resolve_version(v).await?);
resolved = o;
&resolved
}
_ => options,
};
let job = prepare_for(query, params, options, self.caps())?;
if job.writes() && options.query.as_of.is_some() {
return Err(versioned_write());
}
let stats = self.stats.borrow().clone();
let mut sql = SqlCypherJob::new(job, stats, self.caps().clone(), options.query.clone());
let mut response = None;
let r = loop {
let step = match sql.step(response.take()) {
Ok(s) => s,
Err(e) => return Err(sql.take_error().unwrap_or(CypherError::Store(e))),
};
match step {
Step::Execute(request) => response = Some(self.backend.execute(&request).await?),
Step::Done(out) => break out,
}
};
if r.schema_changed {
self.reload_stats().await?;
}
Ok(r)
}
}
fn versioned_write() -> CypherError {
CypherError::Unsupported(
"a writing statement cannot run at a past version: writes apply to the current state"
.into(),
)
}