#![cfg_attr(not(feature = "gql"), allow(dead_code))]
use std::sync::Arc;
use futures::StreamExt;
use crate::exec::{
AccessMode, ContextLevel, ExecOperator, ExecutionContext, FlowResult, OperatorMetrics,
OutputOrdering, ValueBatch, ValueBatchStream, buffer_stream, monitor_stream,
};
use crate::val::{RecordId, Value};
#[derive(Debug, Clone)]
pub struct DistinctEdges {
pub(crate) input: Arc<dyn ExecOperator>,
pub(crate) edge_bindings: Vec<String>,
pub(crate) metrics: Arc<OperatorMetrics>,
}
impl DistinctEdges {
pub(crate) fn new(input: Arc<dyn ExecOperator>, edge_bindings: Vec<String>) -> Self {
Self {
input,
edge_bindings,
metrics: Arc::new(OperatorMetrics::new()),
}
}
}
impl ExecOperator for DistinctEdges {
fn name(&self) -> &'static str {
"DistinctEdges"
}
fn required_context(&self) -> ContextLevel {
ContextLevel::Database.max(self.input.required_context())
}
fn access_mode(&self) -> AccessMode {
self.input.access_mode()
}
fn children(&self) -> Vec<&Arc<dyn ExecOperator>> {
vec![&self.input]
}
fn metrics(&self) -> Option<&OperatorMetrics> {
Some(&self.metrics)
}
fn attrs(&self) -> Vec<(String, String)> {
vec![("edges".to_string(), self.edge_bindings.join(", "))]
}
fn output_ordering(&self) -> OutputOrdering {
self.input.output_ordering()
}
fn execute(&self, ctx: &ExecutionContext) -> FlowResult<ValueBatchStream> {
let input_stream = buffer_stream(
self.input.execute(ctx)?,
self.input.access_mode(),
self.input.cardinality_hint(),
ctx.root().ctx.config.operator_buffer_size,
);
let edge_bindings = self.edge_bindings.clone();
let ctx = ctx.clone();
let filtered = async_stream::try_stream! {
futures::pin_mut!(input_stream);
while let Some(batch_result) = input_stream.next().await {
crate::exec::operators::check_cancelled(&ctx)?;
let batch = batch_result?;
let mut values = Vec::new();
for value in batch.values {
if row_has_distinct_edges(&value, &edge_bindings) {
values.push(value);
}
}
if !values.is_empty() {
yield ValueBatch { values };
}
}
};
Ok(monitor_stream(Box::pin(filtered), "DistinctEdges", &self.metrics))
}
}
fn row_has_distinct_edges(row: &Value, edge_bindings: &[String]) -> bool {
let Value::Object(obj) = row else {
return true;
};
let mut ids: Vec<RecordId> = Vec::new();
for name in edge_bindings {
for rid in slot_edge_ids(obj.get(name.as_str())) {
if !push_distinct(&mut ids, rid) {
return false;
}
}
}
true
}
fn slot_edge_ids(slot: Option<&Value>) -> Vec<RecordId> {
match slot {
Some(Value::Object(edge)) => edge_id_from_object(edge).into_iter().collect(),
Some(Value::RecordId(rid)) => vec![rid.clone()],
Some(Value::Array(group)) => group
.iter()
.filter_map(|elem| match elem {
Value::Object(edge) => edge_id_from_object(edge),
Value::RecordId(rid) => Some(rid.clone()),
_ => None,
})
.collect(),
_ => Vec::new(),
}
}
fn edge_id_from_object(edge: &crate::val::Object) -> Option<RecordId> {
match edge.get("id") {
Some(Value::RecordId(rid)) => Some(rid.clone()),
_ => None,
}
}
fn push_distinct(ids: &mut Vec<RecordId>, rid: RecordId) -> bool {
if ids.contains(&rid) {
return false;
}
ids.push(rid);
true
}
#[cfg(test)]
mod tests {
use super::*;
use crate::exec::operators::test_util::{ValuesOperator, collect, root_ctx};
use crate::val::{Array, Object, RecordId, RecordIdKey, TableName, Value};
fn rid(table: &str, key: &str) -> RecordId {
RecordId {
table: TableName::new(table.to_string()),
key: RecordIdKey::from(key.to_string()),
}
}
fn edge_obj(table: &str, key: &str) -> Value {
let mut o = Object::default();
o.insert("id".to_string(), Value::RecordId(rid(table, key)));
Value::Object(o)
}
fn row(fields: &[(&str, Value)]) -> Value {
let mut o = Object::default();
for (k, v) in fields {
o.insert(k.to_string(), v.clone());
}
Value::Object(o)
}
fn names(ns: &[&str]) -> Vec<String> {
ns.iter().map(|s| s.to_string()).collect()
}
#[tokio::test]
async fn passes_row_with_distinct_edges() {
let r = row(&[("k1", edge_obj("knows", "a")), ("k2", edge_obj("knows", "b"))]);
let input = ValuesOperator::new(vec![r.clone()]);
let op: Arc<dyn ExecOperator> = Arc::new(DistinctEdges::new(input, names(&["k1", "k2"])));
let ctx = root_ctx();
assert_eq!(collect(&op, &ctx).await, vec![r]);
}
#[tokio::test]
async fn drops_row_with_duplicate_edge() {
let r = row(&[("k1", edge_obj("knows", "a")), ("k2", edge_obj("knows", "a"))]);
let input = ValuesOperator::new(vec![r]);
let op: Arc<dyn ExecOperator> = Arc::new(DistinctEdges::new(input, names(&["k1", "k2"])));
let ctx = root_ctx();
assert!(collect(&op, &ctx).await.is_empty());
}
#[tokio::test]
async fn drops_row_with_duplicate_hidden_record_id_edges() {
let r = row(&[
("__e0", Value::RecordId(rid("knows", "x"))),
("__e1", Value::RecordId(rid("knows", "x"))),
]);
let input = ValuesOperator::new(vec![r]);
let op: Arc<dyn ExecOperator> =
Arc::new(DistinctEdges::new(input, names(&["__e0", "__e1"])));
let ctx = root_ctx();
assert!(collect(&op, &ctx).await.is_empty());
}
#[tokio::test]
async fn flattens_group_binding_and_drops_internal_duplicate() {
let group = Value::Array(Array(vec![edge_obj("knows", "a"), edge_obj("knows", "a")]));
let r = row(&[("g", group)]);
let input = ValuesOperator::new(vec![r]);
let op: Arc<dyn ExecOperator> = Arc::new(DistinctEdges::new(input, names(&["g"])));
let ctx = root_ctx();
assert!(collect(&op, &ctx).await.is_empty());
}
#[tokio::test]
async fn passes_group_binding_with_distinct_edges() {
let group = Value::Array(Array(vec![edge_obj("knows", "a"), edge_obj("knows", "b")]));
let r = row(&[("g", group)]);
let input = ValuesOperator::new(vec![r.clone()]);
let op: Arc<dyn ExecOperator> = Arc::new(DistinctEdges::new(input, names(&["g"])));
let ctx = root_ctx();
assert_eq!(collect(&op, &ctx).await, vec![r]);
}
#[tokio::test]
async fn drops_when_single_edge_collides_with_group_element() {
let group = Value::Array(Array(vec![edge_obj("knows", "a"), edge_obj("knows", "b")]));
let r = row(&[("g", group), ("k", edge_obj("knows", "a"))]);
let input = ValuesOperator::new(vec![r]);
let op: Arc<dyn ExecOperator> = Arc::new(DistinctEdges::new(input, names(&["g", "k"])));
let ctx = root_ctx();
assert!(collect(&op, &ctx).await.is_empty());
}
#[tokio::test]
async fn passes_mixed_group_and_distinct_single_edge() {
let group = Value::Array(Array(vec![edge_obj("knows", "a"), edge_obj("knows", "b")]));
let r = row(&[("g", group), ("k", edge_obj("knows", "c"))]);
let input = ValuesOperator::new(vec![r.clone()]);
let op: Arc<dyn ExecOperator> = Arc::new(DistinctEdges::new(input, names(&["g", "k"])));
let ctx = root_ctx();
assert_eq!(collect(&op, &ctx).await, vec![r]);
}
#[tokio::test]
async fn null_binding_is_skipped() {
let r = row(&[("k", edge_obj("knows", "a")), ("opt", Value::Null)]);
let input = ValuesOperator::new(vec![r.clone()]);
let op: Arc<dyn ExecOperator> = Arc::new(DistinctEdges::new(input, names(&["k", "opt"])));
let ctx = root_ctx();
assert_eq!(collect(&op, &ctx).await, vec![r]);
}
#[tokio::test]
async fn empty_group_passes() {
let r = row(&[("g", Value::Array(Array(Vec::new()))), ("k", edge_obj("knows", "a"))]);
let input = ValuesOperator::new(vec![r.clone()]);
let op: Arc<dyn ExecOperator> = Arc::new(DistinctEdges::new(input, names(&["g", "k"])));
let ctx = root_ctx();
assert_eq!(collect(&op, &ctx).await, vec![r]);
}
#[tokio::test]
async fn preserves_order_dropping_only_offenders() {
let keep1 = row(&[("k1", edge_obj("knows", "a")), ("k2", edge_obj("knows", "b"))]);
let drop = row(&[("k1", edge_obj("knows", "c")), ("k2", edge_obj("knows", "c"))]);
let keep2 = row(&[("k1", edge_obj("knows", "d")), ("k2", edge_obj("knows", "e"))]);
let input = ValuesOperator::new(vec![keep1.clone(), drop, keep2.clone()]);
let op: Arc<dyn ExecOperator> = Arc::new(DistinctEdges::new(input, names(&["k1", "k2"])));
let ctx = root_ctx();
assert_eq!(collect(&op, &ctx).await, vec![keep1, keep2]);
}
#[tokio::test]
async fn empty_input_yields_nothing() {
let input = ValuesOperator::new(Vec::new());
let op: Arc<dyn ExecOperator> = Arc::new(DistinctEdges::new(input, names(&["k1", "k2"])));
let ctx = root_ctx();
assert!(collect(&op, &ctx).await.is_empty());
}
#[test]
fn reports_name_and_edges_attr() {
let op = DistinctEdges::new(ValuesOperator::new(Vec::new()), names(&["k1", "k2"]));
assert_eq!(op.name(), "DistinctEdges");
assert_eq!(op.attrs(), vec![("edges".to_string(), "k1, k2".to_string())]);
}
}