use super::helpers::*;
use super::*;
use crate::datatypes::values::Value;
use crate::graph::parallel::{self, ParallelInterrupt};
use crate::graph::storage::GraphRead;
use petgraph::graph::NodeIndex;
use std::collections::{HashMap, HashSet};
fn scoped_node_and_rel(
params: &HashMap<String, Value>,
) -> (
Option<Vec<String>>,
Option<Vec<crate::graph::schema::InternedKey>>,
) {
let node_types = string_list_param(params, "node_type");
let rel_types = string_list_param(params, "relationship").map(|names| {
names
.iter()
.map(|s| crate::graph::schema::InternedKey::from_str(s))
.collect()
});
(node_types, rel_types)
}
fn string_list_param(params: &HashMap<String, Value>, key: &str) -> Option<Vec<String>> {
match params.get(key) {
Some(Value::String(s)) => Some(vec![s.clone()]),
Some(Value::List(items)) => {
let v: Vec<String> = items
.iter()
.filter_map(|x| match x {
Value::String(s) => Some(s.clone()),
_ => None,
})
.collect();
if v.is_empty() {
None
} else {
Some(v)
}
}
_ => None,
}
}
fn parse_scope_predicate(src: &str) -> Result<Predicate, String> {
let wrapped = format!("MATCH (n) WHERE {src} RETURN n");
let query = crate::graph::languages::cypher::parser::parse_cypher(&wrapped)
.map_err(|e| format!("invalid `where` predicate '{src}': {e}"))?;
query
.clauses
.into_iter()
.find_map(|c| match c {
Clause::Where(w) => Some(w.predicate),
_ => None,
})
.ok_or_else(|| format!("`where` predicate '{src}' did not parse to a condition"))
}
fn algo_allowed_keys(proc: &str) -> Option<Vec<&'static str>> {
let mut keys: Vec<&'static str> = match proc {
"pagerank" => vec!["damping_factor", "max_iterations", "tolerance", "where"],
"betweenness" | "betweenness_centrality" => vec!["normalized", "sample_size", "where"],
"closeness" | "closeness_centrality" => vec!["normalized", "sample_size", "where"],
"degree" | "degree_centrality" => vec!["normalized", "where"],
"louvain" | "louvain_communities" | "leiden" | "leiden_communities" => {
vec!["resolution", "weight_property", "where"]
}
"label_propagation" => vec!["max_iterations", "where"],
"ready_set" | "dependency_frontier" => vec!["done"],
"connected_components"
| "weakly_connected_components"
| "k_core"
| "coreness"
| "clustering_coefficient"
| "local_clustering_coefficient"
| "triangle_count"
| "transitivity"
| "eccentricity"
| "diameter" => vec![],
_ => return None,
};
keys.extend([
"node_type",
"node_types",
"relationship",
"connection_types",
"timeout_ms",
]);
Some(keys)
}
const FAIL_OPEN_ON_EMPTY_REL_SCOPE: &[&str] = &["ready_set", "dependency_frontier"];
fn validate_scope_names(
proc: &str,
params: &HashMap<String, Value>,
graph: &crate::graph::DirGraph,
) -> Result<Vec<String>, String> {
if algo_allowed_keys(proc).is_none() {
return Ok(Vec::new());
}
let mut warnings = Vec::new();
if !graph.connection_type_metadata.is_empty() {
let unknown: Vec<String> = string_list_param(params, "relationship")
.unwrap_or_default()
.into_iter()
.filter(|r| !graph.connection_type_metadata.contains_key(r))
.collect();
if !unknown.is_empty() {
let mut valid: Vec<&str> = graph
.connection_type_metadata
.keys()
.map(String::as_str)
.collect();
valid.sort_unstable();
let fatal = FAIL_OPEN_ON_EMPTY_REL_SCOPE.contains(&proc) || graph.schema_locked;
for rel in &unknown {
let hint = crate::graph::mutation::validation::did_you_mean(rel, &valid);
if fatal {
return Err(format!(
"CALL {proc}(): unknown relationship type '{rel}'.{hint}\n \
Valid types: {}\n A relationship type that matches no edges would \
report every node as ready (an empty dependency list satisfies the \
'all dependencies done' test vacuously), so this is refused rather \
than answered.",
valid.join(", ")
));
}
warnings.push(format!(
"CALL {proc}() references unknown relationship type '{rel}' — the graph has \
no such edge type, so the scoped subgraph has no edges of it.{hint}"
));
}
}
}
let have_node_schema =
!graph.node_type_metadata.is_empty() || graph.type_indices.keys().next().is_some();
if have_node_schema {
let unknown: Vec<String> = string_list_param(params, "node_type")
.unwrap_or_default()
.into_iter()
.filter(|t| {
!graph.node_type_metadata.contains_key(t)
&& !graph.type_indices.contains_key(t.as_str())
})
.collect();
if !unknown.is_empty() {
let mut valid: Vec<&str> = graph
.node_type_metadata
.keys()
.map(String::as_str)
.chain(graph.type_indices.keys())
.collect();
valid.sort_unstable();
valid.dedup();
for ty in &unknown {
let hint = crate::graph::mutation::validation::did_you_mean(ty, &valid);
if graph.schema_locked {
return Err(format!(
"CALL {proc}(): unknown node type '{ty}'.{hint}\n Valid types: {}",
valid.join(", ")
));
}
warnings.push(format!(
"CALL {proc}() references unknown node type '{ty}' — the graph has no such \
type, so it contributes no nodes.{hint}"
));
}
}
}
Ok(warnings)
}
fn normalize_and_validate_algo_params(
proc: &str,
params: &mut HashMap<String, Value>,
graph: &crate::graph::DirGraph,
) -> Result<Vec<String>, String> {
let Some(allowed) = algo_allowed_keys(proc) else {
return Ok(Vec::new());
};
fn alias(params: &mut HashMap<String, Value>, from: &str, to: &str) {
if !params.contains_key(to) {
if let Some(v) = params.get(from).cloned() {
params.insert(to.to_string(), v);
}
}
}
alias(params, "relationship", "connection_types");
alias(params, "connection_types", "relationship");
alias(params, "node_types", "node_type");
for key in params.keys() {
if !allowed.contains(&key.as_str()) {
let hint = crate::graph::mutation::validation::did_you_mean(key, &allowed);
return Err(format!("CALL {proc}(): unknown config key '{key}'.{hint}"));
}
}
validate_scope_names(proc, params, graph)
}
impl<'a> CypherExecutor<'a> {
fn validate_algo_params(
&self,
proc: &str,
params: &mut HashMap<String, Value>,
) -> Result<(), String> {
for warning in normalize_and_validate_algo_params(proc, params, self.graph)? {
self.warn(warning);
}
Ok(())
}
fn build_node_scope(
&self,
params: &HashMap<String, Value>,
) -> Result<Option<HashSet<NodeIndex>>, String> {
let node_types = string_list_param(params, "node_type");
let where_src = match params.get("where") {
Some(Value::String(s)) if !s.trim().is_empty() => Some(s.as_str()),
_ => None,
};
if node_types.is_none() && where_src.is_none() {
return Ok(None);
}
let candidates: Vec<NodeIndex> = match &node_types {
Some(types) => {
let mut v = Vec::new();
for t in types {
if let Some(idxs) = self.graph.type_indices.get(t.as_str()) {
v.extend(idxs.iter());
}
}
v
}
None => self.graph.graph.node_indices().collect(),
};
let predicate = match where_src {
Some(src) => Some(parse_scope_predicate(src)?),
None => None,
};
let mut scope = HashSet::with_capacity(candidates.len());
for (i, idx) in candidates.into_iter().enumerate() {
if i & 0xFFFF == 0 {
self.check_deadline()?;
}
if let Some(pred) = &predicate {
let mut row = ResultRow::new();
row.node_bindings.insert("n".to_string(), idx);
if !self.evaluate_predicate(pred, &row)? {
continue;
}
}
scope.insert(idx);
}
Ok(Some(scope))
}
pub(super) fn execute_unwind(
&self,
clause: &UnwindClause,
result_set: ResultSet,
) -> Result<ResultSet, String> {
self.check_deadline()?;
let mut new_rows = Vec::new();
for (row_idx, mut row) in result_set.rows.into_iter().enumerate() {
self.check_interrupt_periodic(row_idx)?;
let val = match (&clause.expression, clause.consume_source) {
(Expression::Variable(name), true) => match row.projected.remove(name.as_str()) {
Some(taken) => taken,
None => self.evaluate_expression(&clause.expression, &row)?,
},
_ => self.evaluate_expression(&clause.expression, &row)?,
};
match val {
Value::List(items) => {
let total = items.len();
self.budget.check_work(total, "UNWIND collection")?;
self.budget.reserve_rows(new_rows.len(), total, "UNWIND")?;
for (i, item_val) in items.into_iter().enumerate() {
self.check_interrupt_periodic(i)?;
if i + 1 == total {
row.projected.insert(clause.alias.clone(), item_val);
new_rows.push(row);
break;
}
let mut new_row = row.clone();
new_row.projected.insert(clause.alias.clone(), item_val);
new_rows.push(new_row);
}
}
Value::String(s) if s.starts_with('[') && s.ends_with(']') => {
let items = split_list_top_level(&s);
let total = items.len();
self.budget.check_work(total, "UNWIND collection")?;
self.budget.reserve_rows(new_rows.len(), total, "UNWIND")?;
for (i, item_str) in items.into_iter().enumerate() {
self.check_interrupt_periodic(i)?;
let parsed_val = parse_value_string(item_str.trim());
if i + 1 == total {
row.projected.insert(clause.alias.clone(), parsed_val);
new_rows.push(row);
break;
}
let mut new_row = row.clone();
new_row.projected.insert(clause.alias.clone(), parsed_val);
new_rows.push(new_row);
}
}
Value::Null => {
}
_ => {
self.budget.reserve_rows(new_rows.len(), 1, "UNWIND")?;
row.projected.insert(clause.alias.clone(), val);
new_rows.push(row);
}
}
}
Ok(ResultSet {
rows: new_rows,
columns: result_set.columns,
lazy_return_items: None,
})
}
pub(super) fn execute_call(
&self,
clause: &CallClause,
existing: ResultSet,
) -> Result<ResultSet, String> {
self.check_deadline()?;
let raw_proc_name = clause.procedure_name.to_lowercase();
let proc_name = raw_proc_name
.strip_prefix("kglite.")
.unwrap_or(raw_proc_name.as_str())
.to_string();
let effective_yields = resolve_yield_items(
proc_name.as_str(),
&clause.procedure_name,
&clause.yield_items,
)?;
let synthesized_clause;
let clause = if clause.yield_items.is_empty() {
synthesized_clause = CallClause {
procedure_name: clause.procedure_name.clone(),
parameters: clause.parameters.clone(),
yield_items: effective_yields,
};
&synthesized_clause
} else {
clause
};
const PROC_FULL_GRAPH_LIMIT: usize = 2_000_000;
let needs_scope = matches!(
proc_name.as_str(),
"pagerank"
| "betweenness"
| "betweenness_centrality"
| "degree"
| "degree_centrality"
| "closeness"
| "closeness_centrality"
| "louvain"
| "louvain_communities"
| "leiden"
| "leiden_communities"
| "label_propagation"
| "connected_components"
| "weakly_connected_components"
);
let streaming_community = matches!(
proc_name.as_str(),
"louvain" | "louvain_communities" | "leiden" | "leiden_communities"
) && (self.graph.graph.is_disk() || self.graph.graph.is_mapped());
let mut params = self.extract_call_params(&clause.parameters)?;
self.validate_algo_params(proc_name.as_str(), &mut params)?;
let scope = if needs_scope {
self.build_node_scope(¶ms)?
} else {
None
};
if needs_scope && self.deadline.is_some() && !streaming_community && scope.is_none() {
let n = self.graph.graph.node_count();
if n > PROC_FULL_GRAPH_LIMIT {
return Err(format!(
"CALL {}() on a graph with {n} nodes would scan the whole graph. \
Scope it with {{node_type: '...', where: '...'}}, try a smaller \
graph, or pass timeout_ms=0 to override this guard.",
clause.procedure_name
));
}
}
let rows = match proc_name.as_str() {
"pagerank"
| "betweenness"
| "betweenness_centrality"
| "degree"
| "degree_centrality"
| "closeness"
| "closeness_centrality"
| "louvain"
| "louvain_communities"
| "leiden"
| "leiden_communities"
| "label_propagation" => super::centrality_procedures::execute_centrality_procedure(
self,
&proc_name,
¶ms,
scope.as_ref(),
streaming_community,
&clause.yield_items,
)?,
"connected_components" | "weakly_connected_components" => {
let (node_types, rel_types) = scoped_node_and_rel(¶ms);
let components =
crate::graph::algorithms::graph_algorithms::weakly_connected_components_scoped(
self.graph,
node_types.as_deref(),
rel_types.as_deref(),
self.interrupt(),
)?;
let mut rows = Vec::new();
let mut row_counter: usize = 0;
for (comp_id, nodes) in components.iter().enumerate() {
for &node_idx in nodes {
row_counter += 1;
if row_counter & 0xFFFFF == 0 {
self.check_deadline()?;
}
let mut row = ResultRow::new();
for item in &clause.yield_items {
let alias = item.alias.as_deref().unwrap_or(&item.name);
match item.name.as_str() {
"node" => {
row.node_bindings.insert(alias.to_string(), node_idx);
}
"component" => {
row.projected
.insert(alias.to_string(), Value::Int64(comp_id as i64));
}
_ => {}
}
}
rows.push(row);
}
}
rows
}
"k_core" | "coreness" => {
let (node_types, rel_types) = scoped_node_and_rel(¶ms);
let scores = crate::graph::algorithms::graph_algorithms::coreness_scoped(
self.graph,
node_types.as_deref(),
rel_types.as_deref(),
self.interrupt(),
)?;
let mut rows = Vec::with_capacity(scores.len());
for (node_idx, core) in scores {
let mut row = ResultRow::new();
for item in &clause.yield_items {
let alias = item.alias.as_deref().unwrap_or(&item.name);
match item.name.as_str() {
"node" => {
row.node_bindings.insert(alias.to_string(), node_idx);
}
"coreness" => {
row.projected.insert(alias.to_string(), Value::Int64(core));
}
_ => {}
}
}
rows.push(row);
}
rows
}
"ready_set" | "dependency_frontier" => {
let (node_types, rel_types) = scoped_node_and_rel(¶ms);
let done_src = match params.get("done") {
Some(Value::String(s)) if !s.trim().is_empty() => s.clone(),
_ => {
return Err("CALL ready_set(): requires a `done` predicate over `n`, \
e.g. done: 'n.status = \"done\"'"
.to_string())
}
};
let predicate = parse_scope_predicate(&done_src)?;
let mut done: HashSet<NodeIndex> = HashSet::new();
for (i, idx) in self.graph.graph.node_indices().enumerate() {
if i & 0xFFFF == 0 {
self.check_deadline()?;
}
let mut row = ResultRow::new();
row.node_bindings.insert("n".to_string(), idx);
if self.evaluate_predicate(&predicate, &row)? {
done.insert(idx);
}
}
let ready = crate::graph::algorithms::graph_algorithms::ready_set_scoped(
self.graph,
node_types.as_deref(),
rel_types.as_deref(),
&done,
self.interrupt(),
)?;
let mut rows = Vec::with_capacity(ready.len());
for (node_idx, dep_count) in ready {
let mut row = ResultRow::new();
for item in &clause.yield_items {
let alias = item.alias.as_deref().unwrap_or(&item.name);
match item.name.as_str() {
"node" => {
row.node_bindings.insert(alias.to_string(), node_idx);
}
"dependency_count" => {
row.projected
.insert(alias.to_string(), Value::Int64(dep_count));
}
_ => {}
}
}
rows.push(row);
}
rows
}
"clustering_coefficient" | "local_clustering_coefficient" => {
let (node_types, rel_types) = scoped_node_and_rel(¶ms);
let scores =
crate::graph::algorithms::graph_algorithms::clustering_coefficient_scoped(
self.graph,
node_types.as_deref(),
rel_types.as_deref(),
self.interrupt(),
)?;
let mut rows = Vec::with_capacity(scores.len());
for (node_idx, coeff) in scores {
let mut row = ResultRow::new();
for item in &clause.yield_items {
let alias = item.alias.as_deref().unwrap_or(&item.name);
match item.name.as_str() {
"node" => {
row.node_bindings.insert(alias.to_string(), node_idx);
}
"coefficient" => {
row.projected
.insert(alias.to_string(), Value::Float64(coeff));
}
_ => {}
}
}
rows.push(row);
}
rows
}
"triangle_count" | "transitivity" => {
let (node_types, rel_types) = scoped_node_and_rel(¶ms);
let (triangles, transitivity) =
crate::graph::algorithms::graph_algorithms::triangle_count_scoped(
self.graph,
node_types.as_deref(),
rel_types.as_deref(),
self.interrupt(),
)?;
let mut row = ResultRow::new();
for item in &clause.yield_items {
let alias = item.alias.as_deref().unwrap_or(&item.name);
match item.name.as_str() {
"triangles" => {
row.projected
.insert(alias.to_string(), Value::Int64(triangles as i64));
}
"transitivity" => {
row.projected
.insert(alias.to_string(), Value::Float64(transitivity));
}
_ => {}
}
}
vec![row]
}
"eccentricity" => {
let (node_types, rel_types) = scoped_node_and_rel(¶ms);
let eccs = crate::graph::algorithms::graph_algorithms::eccentricity_scoped(
self.graph,
node_types.as_deref(),
rel_types.as_deref(),
self.interrupt(),
)?;
let mut rows = Vec::with_capacity(eccs.len());
for (node_idx, ecc) in eccs {
let mut row = ResultRow::new();
for item in &clause.yield_items {
let alias = item.alias.as_deref().unwrap_or(&item.name);
match item.name.as_str() {
"node" => {
row.node_bindings.insert(alias.to_string(), node_idx);
}
"eccentricity" => {
row.projected.insert(alias.to_string(), Value::Int64(ecc));
}
_ => {}
}
}
rows.push(row);
}
rows
}
"diameter" => {
let (node_types, rel_types) = scoped_node_and_rel(¶ms);
let diameter = crate::graph::algorithms::graph_algorithms::diameter_scoped(
self.graph,
node_types.as_deref(),
rel_types.as_deref(),
self.interrupt(),
)?;
let mut row = ResultRow::new();
for item in &clause.yield_items {
let alias = item.alias.as_deref().unwrap_or(&item.name);
if item.name.as_str() == "diameter" {
row.projected
.insert(alias.to_string(), Value::Int64(diameter));
}
}
vec![row]
}
"cluster" => self.execute_call_cluster(¶ms, &clause.yield_items, &existing)?,
"orphan_node"
| "self_loop"
| "cycle_2step"
| "missing_required_edge"
| "missing_inbound_edge"
| "duplicate_title"
| "duplicate_id"
| "outline"
| "null_property"
| "inverse_violation"
| "transitivity_violation"
| "cardinality_violation"
| "type_domain_violation"
| "type_range_violation"
| "parallel_edges"
| "kg_knn" => super::rule_procedures::execute_rule_procedure(
&proc_name,
self.graph,
¶ms,
&clause.yield_items,
)?,
"affected_tests" | "rev_diff" | "dead_code" | "refresh_stats" => {
super::analysis_procedures::execute_analysis_procedure(
&proc_name,
self.graph,
¶ms,
&clause.yield_items,
)?
}
"list_procedures" => {
let mut rows = Vec::new();
for spec in super::procedure_registry::PROCEDURES {
let mut row = ResultRow::new();
for item in &clause.yield_items {
let alias = item.alias.as_deref().unwrap_or(&item.name);
match item.name.as_str() {
"name" => {
row.projected.insert(
alias.to_string(),
Value::String(spec.name.to_string()),
);
}
"description" => {
row.projected.insert(
alias.to_string(),
Value::String(spec.description.to_string()),
);
}
"yield_columns" => {
row.projected.insert(
alias.to_string(),
Value::String(spec.columns.join(", ")),
);
}
_ => {}
}
}
rows.push(row);
}
rows
}
"db.labels" => super::schema_procedures::execute_schema_procedure(
self,
&proc_name,
¶ms,
&clause.yield_items,
)?,
"db.relationshiptypes" => super::schema_procedures::execute_schema_procedure(
self,
&proc_name,
¶ms,
&clause.yield_items,
)?,
"db.indexes" | "db.constraints" => super::schema_procedures::execute_schema_procedure(
self,
&proc_name,
¶ms,
&clause.yield_items,
)?,
"db.propertykeys" => super::schema_procedures::execute_schema_procedure(
self,
&proc_name,
¶ms,
&clause.yield_items,
)?,
"db.schema"
| "db.schema.visualization"
| "db.schema.nodetypeproperties"
| "db.schema.reltypeproperties"
| "apoc.meta.nodetypeproperties"
| "apoc.meta.reltypeproperties" => super::schema_procedures::execute_schema_procedure(
self,
&proc_name,
¶ms,
&clause.yield_items,
)?,
"db.graph_stats" => super::schema_procedures::execute_schema_procedure(
self,
&proc_name,
¶ms,
&clause.yield_items,
)?,
"db.property_stats" => super::schema_procedures::execute_schema_procedure(
self,
&proc_name,
¶ms,
&clause.yield_items,
)?,
"db.property_uniqueness" => super::schema_procedures::execute_schema_procedure(
self,
&proc_name,
¶ms,
&clause.yield_items,
)?,
other if other.starts_with("db.cdc.") => super::cdc_procedures::execute_cdc_procedure(
self,
&proc_name,
¶ms,
&clause.yield_items,
)?,
_ => unreachable!(),
};
self.budget
.check_work(rows.len(), &format!("CALL {proc_name}"))?;
self.budget
.check_rows(rows.len(), &format!("CALL {proc_name}"))?;
Ok(ResultSet {
rows,
columns: clause
.yield_items
.iter()
.map(|item| item.alias.clone().unwrap_or_else(|| item.name.clone()))
.collect(),
lazy_return_items: None,
})
}
pub(super) fn extract_call_params(
&self,
params: &[(String, Expression)],
) -> Result<HashMap<String, Value>, String> {
let empty_row = ResultRow::new();
let mut map = HashMap::new();
for (key, expr) in params {
let val = self.evaluate_expression(expr, &empty_row)?;
map.insert(key.clone(), val);
}
Ok(map)
}
pub(super) fn execute_call_cluster(
&self,
params: &HashMap<String, Value>,
yield_items: &[YieldItem],
existing: &ResultSet,
) -> Result<Vec<ResultRow>, String> {
let method = call_param_opt_string(params, "method")
.unwrap_or_else(|| "dbscan".to_string())
.to_lowercase();
let eps = call_param_f64(params, "eps", 0.5);
let min_points = call_param_usize(params, "min_points", 3);
let k = call_param_usize(params, "k", 5);
let max_iterations = call_param_usize(params, "max_iterations", 100);
let normalize = call_param_bool(params, "normalize", false);
let properties: Option<Vec<String>> = params.get("properties").and_then(|v| {
let items = parse_list_value(v);
if items.is_empty() {
return None;
}
let strs: Vec<String> = items
.into_iter()
.filter_map(|item| match item {
Value::String(s) => Some(s),
_ => None,
})
.collect();
if strs.is_empty() {
None
} else {
Some(strs)
}
});
let mut node_indices: Vec<NodeIndex> = Vec::new();
let mut seen: HashSet<NodeIndex> = HashSet::new();
for (row_idx, row) in existing.rows.iter().enumerate() {
self.check_interrupt_periodic(row_idx)?;
for (_, &idx) in row.node_bindings.iter() {
if seen.insert(idx) {
node_indices.push(idx);
}
}
}
if node_indices.is_empty() {
return Err("cluster() requires a preceding MATCH clause that binds nodes".to_string());
}
if method != "dbscan" && method != "kmeans" {
return Err(format!(
"Unknown clustering method '{}'. Available: dbscan, kmeans",
method
));
}
let assignments = if let Some(ref prop_names) = properties {
let mut features: Vec<Vec<f64>> = Vec::new();
let mut valid_indices: Vec<usize> = Vec::new();
for (i, &idx) in node_indices.iter().enumerate() {
self.check_interrupt_periodic(i)?;
if let Some(node) = self.graph.graph.node_view(idx) {
let mut vals = Vec::with_capacity(prop_names.len());
let mut all_present = true;
for prop in prop_names {
if let Some(val) = node.get_property(prop) {
if let Some(f) = value_to_f64(&val) {
vals.push(f);
} else {
all_present = false;
break;
}
} else {
all_present = false;
break;
}
}
if all_present {
features.push(vals);
valid_indices.push(i);
}
}
}
if features.is_empty() {
return Err(format!(
"No nodes have all required numeric properties: {:?}",
prop_names
));
}
if normalize {
crate::graph::algorithms::clustering::normalize_features(&mut features);
}
let cluster_assignments = match method.as_str() {
"dbscan" => {
let dm = crate::graph::algorithms::clustering::euclidean_distance_matrix(
&features,
self.interrupt(),
);
self.check_deadline()?;
crate::graph::algorithms::clustering::dbscan(
&dm,
eps,
min_points,
self.interrupt(),
)
}
"kmeans" => crate::graph::algorithms::clustering::kmeans(
&features,
k,
max_iterations,
self.interrupt(),
),
_ => unreachable!(),
};
cluster_assignments
.into_iter()
.map(|ca| (node_indices[valid_indices[ca.index]], ca.cluster))
.collect::<Vec<_>>()
} else {
let mut points: Vec<(f64, f64)> = Vec::new();
let mut valid_indices: Vec<usize> = Vec::new();
for (i, &idx) in node_indices.iter().enumerate() {
self.check_interrupt_periodic(i)?;
if let Some(node) = self.graph.graph.node_view(idx) {
if let Some(config) = self
.graph
.get_spatial_config(node.node_type_str(&self.graph.interner))
{
let (lat_f, lon_f) = config
.location
.as_ref()
.map(|(a, b)| (a.as_str(), b.as_str()))
.unwrap_or(("latitude", "longitude"));
let geom_fallback = config.geometry.as_deref();
if let Some((lat, lon)) = crate::graph::features::spatial::node_location(
node,
lat_f,
lon_f,
geom_fallback,
) {
points.push((lat, lon));
valid_indices.push(i);
}
}
}
}
if points.is_empty() {
return Err(
"No nodes have spatial data. Either configure spatial fields with \
set_spatial_config() or provide explicit 'properties' parameter."
.to_string(),
);
}
let cluster_assignments = match method.as_str() {
"dbscan" => {
let dm = crate::graph::algorithms::clustering::haversine_distance_matrix(
&points,
self.interrupt(),
);
self.check_deadline()?;
crate::graph::algorithms::clustering::dbscan(
&dm,
eps,
min_points,
self.interrupt(),
)
}
"kmeans" => {
let features: Vec<Vec<f64>> =
points.iter().map(|(lat, lon)| vec![*lat, *lon]).collect();
crate::graph::algorithms::clustering::kmeans(
&features,
k,
max_iterations,
self.interrupt(),
)
}
_ => unreachable!(),
};
cluster_assignments
.into_iter()
.map(|ca| (node_indices[valid_indices[ca.index]], ca.cluster))
.collect::<Vec<_>>()
};
let mut rows = Vec::with_capacity(assignments.len());
self.check_deadline()?;
for (row_idx, (node_idx, cluster_id)) in assignments.iter().enumerate() {
self.check_interrupt_periodic(row_idx)?;
let mut row = ResultRow::new();
for item in yield_items {
let alias = item.alias.as_deref().unwrap_or(&item.name);
match item.name.as_str() {
"node" => {
row.node_bindings.insert(alias.to_string(), *node_idx);
}
"cluster" => {
row.projected
.insert(alias.to_string(), Value::Int64(*cluster_id));
}
_ => {}
}
}
rows.push(row);
}
Ok(rows)
}
pub(super) fn centrality_to_rows(
&self,
results: &[crate::graph::algorithms::graph_algorithms::CentralityResult],
yield_items: &[YieldItem],
) -> Result<Vec<ResultRow>, String> {
let mut rows = Vec::with_capacity(results.len());
for (i, cr) in results.iter().enumerate() {
self.check_interrupt_periodic(i)?;
let mut row = ResultRow::new();
for item in yield_items {
let alias = item.alias.as_deref().unwrap_or(&item.name);
match item.name.as_str() {
"node" => {
row.node_bindings.insert(alias.to_string(), cr.node_idx);
}
"score" => {
row.projected
.insert(alias.to_string(), Value::Float64(cr.score));
}
_ => {}
}
}
rows.push(row);
}
Ok(rows)
}
pub(super) fn community_result_to_rows(
&self,
result: &crate::graph::algorithms::graph_algorithms::CommunityResult,
yield_items: &[YieldItem],
) -> Result<Vec<ResultRow>, String> {
let wants_level = yield_items.iter().any(|y| y.name == "level");
let levels: Vec<&[crate::graph::algorithms::graph_algorithms::CommunityAssignment]> =
if wants_level && !result.levels.is_empty() {
result.levels.iter().map(|v| v.as_slice()).collect()
} else {
vec![result.assignments.as_slice()]
};
let mut rows = Vec::new();
let mut counter = 0usize;
for (lvl, assignments) in levels.iter().enumerate() {
for ca in assignments.iter() {
self.check_interrupt_periodic(counter)?;
counter = counter.saturating_add(1);
let mut row = ResultRow::new();
for item in yield_items {
let alias = item.alias.as_deref().unwrap_or(&item.name);
match item.name.as_str() {
"node" => {
row.node_bindings.insert(alias.to_string(), ca.node_idx);
}
"community" => {
row.projected
.insert(alias.to_string(), Value::Int64(ca.community_id as i64));
}
"level" => {
row.projected
.insert(alias.to_string(), Value::Int64(lvl as i64));
}
_ => {}
}
}
rows.push(row);
}
}
Ok(rows)
}
pub(super) fn execute_union(
&self,
clause: &UnionClause,
result_set: ResultSet,
) -> Result<ResultSet, String> {
let right_result = self.execute(&clause.query)?;
if !result_set.columns.is_empty() && result_set.columns != right_result.columns {
let op = match clause.kind {
SetOpKind::Union => "UNION",
SetOpKind::Intersect => "INTERSECT",
SetOpKind::Except => "EXCEPT",
};
return Err(format!(
"All sub queries in a {op} must have the same return column names \
(left side {:?} != right side {:?}).",
result_set.columns, right_result.columns,
));
}
let columns = if result_set.columns.is_empty() {
right_result.columns.clone()
} else {
result_set.columns.clone()
};
let row_hash = |row: &ResultRow, cols: &[String]| -> u64 {
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
for col in cols {
match row.projected.get(col) {
Some(val) => val.hash(&mut hasher),
None => Value::Null.hash(&mut hasher),
}
}
hasher.finish()
};
match clause.kind {
SetOpKind::Union => {
let mut combined_rows = result_set.rows;
self.budget
.reserve_rows(combined_rows.len(), right_result.rows.len(), "UNION")?;
for (row_idx, row_values) in right_result.rows.into_iter().enumerate() {
self.check_interrupt_periodic(row_idx)?;
let mut projected = Bindings::with_capacity(right_result.columns.len());
for (i, col) in right_result.columns.iter().enumerate() {
if let Some(val) = row_values.get(i) {
projected.insert(col.clone(), val.clone());
}
}
combined_rows.push(ResultRow::from_projected(projected));
}
if !clause.all {
let mut seen = HashSet::new();
let mut deduplicated = Vec::with_capacity(combined_rows.len());
for (row_idx, row) in combined_rows.into_iter().enumerate() {
self.check_interrupt_periodic(row_idx)?;
if seen.insert(row_hash(&row, &columns)) {
deduplicated.push(row);
}
}
combined_rows = deduplicated;
}
Ok(ResultSet {
rows: combined_rows,
columns,
lazy_return_items: None,
})
}
SetOpKind::Intersect => {
self.budget
.consume_collection(right_result.rows.len(), "INTERSECT right-side hash set")?;
let right_columns = right_result.columns.clone();
let mut right_hashes = HashSet::with_capacity(right_result.rows.len());
for (row_idx, row_values) in right_result.rows.iter().enumerate() {
self.check_interrupt_periodic(row_idx)?;
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
for (i, col) in columns.iter().enumerate() {
let val = right_columns
.iter()
.position(|rc| rc == col)
.and_then(|pos| row_values.get(pos))
.or_else(|| row_values.get(i));
match val {
Some(v) => v.hash(&mut hasher),
None => Value::Null.hash(&mut hasher),
}
}
right_hashes.insert(hasher.finish());
}
let mut seen = HashSet::new();
let mut kept = Vec::new();
for (row_idx, row) in result_set.rows.into_iter().enumerate() {
self.check_interrupt_periodic(row_idx)?;
let h = row_hash(&row, &columns);
if right_hashes.contains(&h) && seen.insert(h) {
kept.push(row);
}
}
Ok(ResultSet {
rows: kept,
columns,
lazy_return_items: None,
})
}
SetOpKind::Except => {
self.budget
.consume_collection(right_result.rows.len(), "EXCEPT right-side hash set")?;
let right_columns = right_result.columns.clone();
let mut right_hashes = HashSet::with_capacity(right_result.rows.len());
for (row_idx, row_values) in right_result.rows.iter().enumerate() {
self.check_interrupt_periodic(row_idx)?;
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
for (i, col) in columns.iter().enumerate() {
let val = right_columns
.iter()
.position(|rc| rc == col)
.and_then(|pos| row_values.get(pos))
.or_else(|| row_values.get(i));
match val {
Some(v) => v.hash(&mut hasher),
None => Value::Null.hash(&mut hasher),
}
}
right_hashes.insert(hasher.finish());
}
let mut seen = HashSet::new();
let mut kept = Vec::new();
for (row_idx, row) in result_set.rows.into_iter().enumerate() {
self.check_interrupt_periodic(row_idx)?;
let h = row_hash(&row, &columns);
if !right_hashes.contains(&h) && seen.insert(h) {
kept.push(row);
}
}
Ok(ResultSet {
rows: kept,
columns,
lazy_return_items: None,
})
}
}
}
pub fn finalize_result(&self, mut result_set: ResultSet) -> Result<CypherResult, String> {
if result_set.columns.is_empty() {
if result_set.rows.is_empty() {
return Ok(CypherResult::empty());
}
let first_row = &result_set.rows[0];
let mut columns = Vec::new();
for name in first_row.node_bindings.keys() {
columns.push(name.clone());
}
for name in first_row.edge_bindings.keys() {
columns.push(name.clone());
}
for name in first_row.projected.keys() {
columns.push(name.clone());
}
columns.sort();
let rows: Vec<Vec<Value>> = result_set
.rows
.iter()
.map(|row| {
columns
.iter()
.map(|col| {
if let Some(val) = row.projected.get(col) {
val.clone()
} else if let Some(&idx) = row.node_bindings.get(col) {
if let Some(node) = self.graph.graph.node_view(idx) {
node_to_map_value(node)
} else {
Value::Null
}
} else {
Value::Null
}
})
.collect()
})
.collect();
return Ok(CypherResult {
columns,
rows,
stats: None,
profile: None,
diagnostics: None,
lazy: None,
});
}
if let Some(return_items) = result_set.lazy_return_items.take() {
return Ok(CypherResult {
columns: result_set.columns,
rows: Vec::new(),
stats: None,
profile: None,
diagnostics: None,
lazy: Some(super::super::result::LazyResultDescriptor::new(
result_set.rows,
return_items,
self.graph,
)),
});
}
let columns = std::mem::take(&mut result_set.columns);
let rows: Vec<Vec<Value>> = if result_set.rows.len() >= parallel::PROJECTION_MIN_ROWS {
let cols = &columns;
let interrupt = ParallelInterrupt::new(|| self.check_deadline().err());
let src = &mut result_set.rows;
parallel::install(|| {
src.par_iter_mut()
.enumerate()
.map(|(i, row)| {
interrupt.check(i)?;
Ok(cols
.iter()
.map(|col| row.projected.remove(col).unwrap_or(Value::Null))
.collect())
})
.collect::<Result<Vec<Vec<Value>>, String>>()
})?
} else {
let cols = &columns;
result_set
.rows
.into_iter()
.map(|mut row| {
cols.iter()
.map(|col| row.projected.remove(col).unwrap_or(Value::Null))
.collect()
})
.collect()
};
Ok(CypherResult {
columns,
rows,
stats: None,
profile: None,
diagnostics: None,
lazy: None,
})
}
}
pub(super) fn names_to_rows(names: &[String], yield_items: &[YieldItem]) -> Vec<ResultRow> {
let mut rows = Vec::with_capacity(names.len());
for name in names {
let mut row = ResultRow::new();
for item in yield_items {
let alias = item.alias.as_deref().unwrap_or(&item.name);
row.projected
.insert(alias.to_string(), Value::String(name.clone()));
}
rows.push(row);
}
rows
}
pub(super) fn resolve_yield_items(
proc_name: &str,
display_name: &str,
requested: &[YieldItem],
) -> Result<Vec<YieldItem>, String> {
let valid = valid_yield_columns(proc_name, display_name)?;
if requested.is_empty() {
return Ok(valid
.iter()
.map(|name| YieldItem {
name: (*name).to_string(),
alias: None,
})
.collect());
}
for item in requested {
if !valid.contains(&item.name.as_str()) {
return Err(format!(
"Procedure '{}' does not yield '{}'. Available: {}",
display_name,
item.name,
valid.join(", ")
));
}
}
Ok(requested.to_vec())
}
fn valid_yield_columns(
proc_name: &str,
display_name: &str,
) -> Result<&'static [&'static str], String> {
match super::procedure_registry::find_procedure(proc_name) {
Some(spec) => Ok(spec.columns),
None => Err(format!(
"Unknown procedure '{}'. Available: {}",
display_name,
super::procedure_registry::PROCEDURES
.iter()
.map(|spec| spec.name)
.collect::<Vec<_>>()
.join(", ")
)),
}
}
pub(super) fn indexes_to_rows(
infos: &[crate::graph::introspection::schema_overview::IndexInfo],
yield_items: &[YieldItem],
) -> Vec<ResultRow> {
let mut rows = Vec::with_capacity(infos.len());
for info in infos {
let mut row = ResultRow::new();
for item in yield_items {
let alias = item.alias.as_deref().unwrap_or(&item.name);
let val = match item.name.as_str() {
"name" => Value::String(info.name.clone()),
"type" => Value::String(info.kind.neo4j_type().to_string()),
"entityType" => Value::String(info.entity_type.to_string()),
"labelsOrTypes" => Value::List(
info.labels_or_types
.iter()
.cloned()
.map(Value::String)
.collect(),
),
"properties" => {
Value::List(info.properties.iter().cloned().map(Value::String).collect())
}
"state" => Value::String(info.state.to_string()),
_ => continue, };
row.projected.insert(alias.to_string(), val);
}
rows.push(row);
}
rows
}
pub(super) fn constraints_to_rows(
infos: &[crate::graph::introspection::schema_overview::ConstraintInfo],
yield_items: &[YieldItem],
) -> Vec<ResultRow> {
let mut rows = Vec::with_capacity(infos.len());
for info in infos {
let mut row = ResultRow::new();
for item in yield_items {
let alias = item.alias.as_deref().unwrap_or(&item.name);
let val = match item.name.as_str() {
"name" => Value::String(info.name.clone()),
"type" => Value::String(info.neo4j_type().to_string()),
"entityType" => Value::String(info.entity_type().to_string()),
"labelsOrTypes" => Value::List(
info.labels_or_types
.iter()
.cloned()
.map(Value::String)
.collect(),
),
"properties" => {
Value::List(info.properties.iter().cloned().map(Value::String).collect())
}
"propertyType" => info
.property_type
.map(|declared| Value::String(declared.name().to_string()))
.unwrap_or(Value::Null),
_ => continue, };
row.projected.insert(alias.to_string(), val);
}
rows.push(row);
}
rows
}
pub(super) fn compute_property_stats(
executor: &CypherExecutor<'_>,
node_type: &str,
prop_name: &str,
) -> Result<(i64, i64, i64), String> {
use std::collections::HashSet;
let graph = executor.graph;
let Some(indices) = graph.type_indices.get(node_type) else {
return Ok((0, 0, 0));
};
let mut value_count: i64 = 0;
let mut null_count: i64 = 0;
let mut seen = HashSet::new();
for (node_count, node_idx) in indices.iter().enumerate() {
executor.check_interrupt_periodic(node_count)?;
let Some(node) = graph.graph.node_view(node_idx) else {
continue;
};
match node.get_field_ref(prop_name) {
Some(v) if !matches!(*v, crate::datatypes::values::Value::Null) => {
value_count += 1;
seen.insert(format!("{v:?}"));
}
_ => {
null_count += 1;
}
}
}
Ok((value_count, null_count, seen.len() as i64))
}
#[cfg(test)]
mod scope_name_tests {
use super::*;
use crate::graph::DirGraph;
fn graph_with_schema() -> DirGraph {
let mut g = DirGraph::new();
g.upsert_node_type_metadata("Task", HashMap::new());
g.upsert_node_type_metadata("Spec", HashMap::new());
g.upsert_connection_type_metadata("DEPENDS_ON", "Task", "Task", HashMap::new());
g.upsert_connection_type_metadata("IMPLEMENTS", "Task", "Spec", HashMap::new());
g
}
fn params(pairs: &[(&str, &str)]) -> HashMap<String, Value> {
pairs
.iter()
.map(|(k, v)| ((*k).to_string(), Value::String((*v).to_string())))
.collect()
}
#[test]
fn ready_set_refuses_an_unknown_relationship_type() {
let g = graph_with_schema();
let error =
validate_scope_names("ready_set", ¶ms(&[("relationship", "DEPENDS_O")]), &g)
.expect_err("a bogus dependency edge must not report every node ready");
assert!(
error.contains("unknown relationship type 'DEPENDS_O'"),
"{error}"
);
assert!(
error.contains("Did you mean 'DEPENDS_ON'"),
"a one-character typo must get the suggestion: {error}"
);
assert!(
error.contains("DEPENDS_ON, IMPLEMENTS"),
"the valid set is what makes the error actionable: {error}"
);
assert!(
validate_scope_names(
"dependency_frontier",
¶ms(&[("relationship", "NOPE")]),
&g
)
.is_err(),
"both spellings of the fail-open procedure must refuse"
);
}
#[test]
fn other_procedures_warn_instead_of_failing() {
let g = graph_with_schema();
for proc in [
"pagerank",
"connected_components",
"k_core",
"triangle_count",
] {
let warnings = validate_scope_names(proc, ¶ms(&[("relationship", "NOPE")]), &g)
.unwrap_or_else(|e| panic!("{proc} must warn, not refuse: {e}"));
assert_eq!(warnings.len(), 1, "{proc}: {warnings:?}");
assert!(
warnings[0].contains("unknown relationship type 'NOPE'"),
"{warnings:?}"
);
}
}
#[test]
fn a_known_relationship_type_is_silent() {
let g = graph_with_schema();
for proc in ["ready_set", "pagerank"] {
assert!(
validate_scope_names(proc, ¶ms(&[("relationship", "DEPENDS_ON")]), &g)
.expect("a known type must pass")
.is_empty(),
"{proc} warned about a type the graph has"
);
}
}
#[test]
fn an_unknown_node_type_warns_and_a_locked_schema_refuses() {
let mut g = graph_with_schema();
let p = params(&[("node_type", "Tsk")]);
let warnings = validate_scope_names("pagerank", &p, &g).expect("open schema warns");
assert_eq!(warnings.len(), 1, "{warnings:?}");
assert!(warnings[0].contains("Did you mean 'Task'"), "{warnings:?}");
g.schema_locked = true;
let error = validate_scope_names("pagerank", &p, &g)
.expect_err("a locked schema declares its type set final");
assert!(error.contains("unknown node type 'Tsk'"), "{error}");
assert!(
validate_scope_names("pagerank", ¶ms(&[("relationship", "NOPE")]), &g).is_err(),
"a locked schema refuses an unknown relationship type on every procedure"
);
}
#[test]
fn a_graph_without_edge_metadata_is_not_second_guessed() {
let empty = DirGraph::new();
assert!(
validate_scope_names("ready_set", ¶ms(&[("relationship", "ANY")]), &empty)
.expect("an edgeless graph must not refuse")
.is_empty()
);
}
#[test]
fn non_algorithm_procedures_are_untouched() {
let g = graph_with_schema();
assert!(
validate_scope_names("orphan_node", ¶ms(&[("relationship", "NOPE")]), &g)
.expect("no validation for non-algorithm procedures")
.is_empty()
);
}
}