use anyhow::Result;
use instant::Instant;
use std::collections::HashMap;
use velesdb_core::Database;
use crate::repl::{QueryKind, QueryResult};
use crate::session::SessionSettings;
pub fn execute_query(
db: &Database,
query: &str,
active_collection: Option<&str>,
session: Option<&SessionSettings>,
) -> Result<QueryResult> {
let start = Instant::now();
let mut parsed = velesdb_core::velesql::Parser::parse(query)
.map_err(|e| anyhow::anyhow!("Parse error: {}", e.message))?;
if has_param_vector(&parsed) {
return Err(param_vector_unsupported_error());
}
if let Some(session) = session {
apply_session_settings(&mut parsed, session);
}
if !parsed.is_match_query() && parsed.select.is_aggregation_query() {
return run_aggregation_query(db, &parsed, active_collection, start);
}
run_row_query(db, &parsed, active_collection, start)
}
pub fn apply_session_settings(
parsed: &mut velesdb_core::velesql::Query,
session: &SessionSettings,
) {
if parsed.is_match_query() {
return;
}
inject_session_with_options(&mut parsed.select, session);
cap_limit_to_max_results(&mut parsed.select, session.max_results());
}
fn inject_session_with_options(
select: &mut velesdb_core::velesql::SelectStatement,
session: &SessionSettings,
) {
use velesdb_core::velesql::{WithClause, WithValue};
let with = select.with_clause.get_or_insert_with(WithClause::new);
if with.get_mode().is_none() {
with.options.push(velesdb_core::velesql::WithOption {
key: "mode".to_string(),
value: WithValue::String(session.mode_str()),
});
}
if let Some(ef) = session.ef_search() {
if with.get_ef_search().is_none() {
with.options.push(velesdb_core::velesql::WithOption {
key: "ef_search".to_string(),
value: WithValue::Integer(i64::try_from(ef).unwrap_or(i64::MAX)),
});
}
}
}
fn cap_limit_to_max_results(
select: &mut velesdb_core::velesql::SelectStatement,
max_results: usize,
) {
let cap = max_results as u64;
select.limit = Some(select.limit.map_or(cap, |l| l.min(cap)));
}
fn has_param_vector(parsed: &velesdb_core::velesql::Query) -> bool {
parsed
.select
.where_clause
.as_ref()
.is_some_and(contains_param_vector)
|| parsed
.match_clause
.as_ref()
.and_then(|m| m.where_clause.as_ref())
.is_some_and(contains_param_vector)
}
fn param_vector_unsupported_error() -> anyhow::Error {
anyhow::anyhow!(
"Vector search with $parameter requires the REST API. \
Use literal vectors or metadata-only queries."
)
}
fn run_row_query(
db: &Database,
parsed: &velesdb_core::velesql::Query,
active_collection: Option<&str>,
start: Instant,
) -> Result<QueryResult> {
let kind = query_kind(parsed);
let rows = if parsed.is_match_query() {
let results = route_match_query(db, parsed, active_collection)?;
results.into_iter().map(result_to_row).collect()
} else {
let results = db
.execute_query(parsed, &HashMap::new())
.map_err(|e| anyhow::anyhow!("Query error: {e}"))?;
if matches!(kind, QueryKind::Select) {
project_select_rows(&results, &parsed.select.columns)
} else {
results.into_iter().map(result_to_row).collect()
}
};
Ok(QueryResult {
rows,
duration_ms: start.elapsed().as_secs_f64() * 1000.0,
kind,
})
}
fn project_select_rows(
results: &[velesdb_core::SearchResult],
columns: &velesdb_core::velesql::SelectColumns,
) -> Vec<HashMap<String, serde_json::Value>> {
velesdb_core::collection::search::query::projection::project_results(results, columns)
.into_iter()
.map(json_object_to_row)
.collect()
}
fn query_kind(parsed: &velesdb_core::velesql::Query) -> QueryKind {
if parsed.is_ddl_query() {
QueryKind::Ddl
} else if parsed.is_introspection_query() {
QueryKind::Introspection
} else if parsed.is_admin_query() {
QueryKind::Admin
} else if parsed.is_dml_query() {
QueryKind::Dml
} else if parsed.is_train() {
QueryKind::Train
} else {
QueryKind::Select
}
}
fn route_match_query(
db: &Database,
parsed: &velesdb_core::velesql::Query,
active_collection: Option<&str>,
) -> Result<Vec<velesdb_core::SearchResult>> {
let params = params_with_active_collection(
parsed,
active_collection,
"MATCH queries require an active collection. Use: .use <collection_name>",
)?;
db.execute_query(parsed, ¶ms)
.map_err(|e| anyhow::anyhow!("Query error: {e}"))
}
fn params_with_active_collection(
parsed: &velesdb_core::velesql::Query,
active_collection: Option<&str>,
requires_msg: &str,
) -> Result<HashMap<String, serde_json::Value>> {
let mut params = HashMap::new();
if parsed.select.from.is_empty() {
let col_name = active_collection.ok_or_else(|| anyhow::anyhow!("{requires_msg}"))?;
params.insert(
"_collection".to_string(),
serde_json::Value::String(col_name.to_string()),
);
}
Ok(params)
}
fn run_aggregation_query(
db: &Database,
parsed: &velesdb_core::velesql::Query,
active_collection: Option<&str>,
start: Instant,
) -> Result<QueryResult> {
let params = params_with_active_collection(
parsed,
active_collection,
"Aggregation queries require an active collection. Use: .use <collection_name>",
)?;
let value = db
.execute_aggregate(parsed, ¶ms)
.map_err(|e| anyhow::anyhow!("Query error: {e}"))?;
Ok(QueryResult {
rows: aggregate_value_to_rows(value),
duration_ms: start.elapsed().as_secs_f64() * 1000.0,
kind: QueryKind::Select,
})
}
fn aggregate_value_to_rows(value: serde_json::Value) -> Vec<HashMap<String, serde_json::Value>> {
match value {
serde_json::Value::Array(items) => items.into_iter().map(json_object_to_row).collect(),
other => vec![json_object_to_row(other)],
}
}
fn json_object_to_row(value: serde_json::Value) -> HashMap<String, serde_json::Value> {
match value {
serde_json::Value::Object(map) => map.into_iter().collect(),
other => HashMap::from([("value".to_string(), other)]),
}
}
fn result_to_row(r: velesdb_core::SearchResult) -> HashMap<String, serde_json::Value> {
let mut row = HashMap::new();
row.insert("id".to_string(), serde_json::json!(r.point.id));
row.insert("score".to_string(), serde_json::json!(r.score));
if let Some(serde_json::Value::Object(map)) = &r.point.payload {
for (k, v) in map {
row.insert(k.clone(), v.clone());
}
}
row
}
pub(crate) fn contains_param_vector(condition: &velesdb_core::velesql::Condition) -> bool {
use velesdb_core::velesql::{Condition, SparseVectorExpr, VectorExpr};
match condition {
Condition::VectorSearch(vs) => matches!(vs.vector, VectorExpr::Parameter(_)),
Condition::VectorFusedSearch(vfs) => vfs
.vectors
.iter()
.any(|v| matches!(v, VectorExpr::Parameter(_))),
Condition::SparseVectorSearch(svs) => {
matches!(svs.vector, SparseVectorExpr::Parameter(_))
}
Condition::Similarity(sim) => matches!(sim.vector, VectorExpr::Parameter(_)),
Condition::And(left, right) | Condition::Or(left, right) => {
contains_param_vector(left) || contains_param_vector(right)
}
Condition::Not(inner) | Condition::Group(inner) => contains_param_vector(inner),
Condition::Comparison(_)
| Condition::In(_)
| Condition::Between(_)
| Condition::Like(_)
| Condition::IsNull(_)
| Condition::Match(_)
| Condition::GraphMatch(_)
| Condition::Contains(_)
| Condition::ContainsText(_)
| Condition::GeoDistance(_)
| Condition::GeoBbox(_) => false,
_ => false,
}
}
#[cfg(test)]
#[path = "repl_execute_tests.rs"]
mod repl_execute_tests;