Skip to main content

graphrecords_query/operations/traversal/
nodes.rs

1use crate::{
2    Definite, EntityReference, EvaluateOperand, Explain, IndexDomain, Indexed, Multiple, Operand,
3    OrderState, QueryResult, Single, Unit, Unordered,
4    execution::EvaluationCache,
5    operands::NodesOperand,
6    operations::{Apply, KeyedStream, LaneKernel, Operation, OperationContext, Prepare},
7    optimizer::{OperationInputs, OptimizerHints, PlanIdentity, PlanInputs},
8    registry::operation_manifest,
9    traits::Nodes,
10};
11use graphrecords_core::{GraphRecord, graphrecord::EdgeIndex};
12use graphrecords_utils::aliases::GrHashSet;
13use std::iter::{empty, once};
14
15#[derive(Clone, Explain, Operation, OperationInputs, OptimizerHints, PlanIdentity, PlanInputs)]
16#[operation(scope = Lane)]
17#[explain(label = "Nodes")]
18#[plan(optimizer_hints(empty = if_any))]
19pub struct NodesOperation;
20
21impl Prepare for NodesOperation {
22    type Prepared<'a> = ();
23
24    fn prepare<'a>(
25        &'a self,
26        _graphrecord: &'a GraphRecord,
27        _cache: &'a EvaluationCache<'a>,
28    ) -> QueryResult<Self::Prepared<'a>> {
29        Ok(())
30    }
31}
32
33impl<O: OrderState> LaneKernel<Indexed<EdgeIndex, Unit>, Multiple<O>> for NodesOperation {
34    type Output = NodesOperand<Unordered>;
35
36    fn execute<'a>(
37        graphrecord: &'a GraphRecord,
38        values: KeyedStream<'a, EdgeIndex, Unit, Multiple<O>>,
39        _prepared: Self::Prepared<'a>,
40    ) -> QueryResult<<Self::Output as EvaluateOperand>::ReturnValue<'a>> {
41        let mut nodes = GrHashSet::default();
42
43        for (edge, membership) in values {
44            membership?;
45            let (source, target) = graphrecord.edge_endpoints(edge).expect("Edge must exist");
46            nodes.insert(source);
47            nodes.insert(target);
48        }
49
50        Ok(Box::new(nodes.into_iter().map(|node| (node, Ok(())))))
51    }
52}
53
54impl LaneKernel<Indexed<EdgeIndex, Unit>, Single> for NodesOperation {
55    type Output = NodesOperand<Unordered>;
56
57    fn execute<'a>(
58        graphrecord: &'a GraphRecord,
59        value: KeyedStream<'a, EdgeIndex, Unit, Single>,
60        _prepared: Self::Prepared<'a>,
61    ) -> QueryResult<<Self::Output as EvaluateOperand>::ReturnValue<'a>> {
62        let Some((edge, membership)) = value else {
63            return Ok(Box::new(empty()));
64        };
65        membership?;
66
67        let (source, target) = graphrecord.edge_endpoints(edge).expect("Edge must exist");
68        let nodes: GrHashSet<_> = once(source).chain(once(target)).collect();
69
70        Ok(Box::new(nodes.into_iter().map(|node| (node, Ok(())))))
71    }
72}
73
74impl LaneKernel<Indexed<EdgeIndex, Unit>, Definite> for NodesOperation {
75    type Output = NodesOperand<Unordered>;
76
77    fn execute<'a>(
78        graphrecord: &'a GraphRecord,
79        value: KeyedStream<'a, EdgeIndex, Unit, Definite>,
80        _prepared: Self::Prepared<'a>,
81    ) -> QueryResult<<Self::Output as EvaluateOperand>::ReturnValue<'a>> {
82        let (edge, membership) = value;
83        membership?;
84
85        let (source, target) = graphrecord.edge_endpoints(edge).expect("Edge must exist");
86        let nodes: GrHashSet<_> = once(source).chain(once(target)).collect();
87
88        Ok(Box::new(nodes.into_iter().map(|node| (node, Ok(())))))
89    }
90}
91
92impl<I: IndexDomain, O: OrderState> LaneKernel<Indexed<I, EntityReference<EdgeIndex>>, Multiple<O>>
93    for NodesOperation
94{
95    type Output = NodesOperand<Unordered>;
96
97    fn execute<'a>(
98        graphrecord: &'a GraphRecord,
99        values: KeyedStream<'a, I, EntityReference<EdgeIndex>, Multiple<O>>,
100        _prepared: Self::Prepared<'a>,
101    ) -> QueryResult<<Self::Output as EvaluateOperand>::ReturnValue<'a>> {
102        let mut nodes = GrHashSet::default();
103
104        for value in values {
105            let edge = value.1?;
106            let (source, target) = graphrecord.edge_endpoints(edge).expect("Edge must exist");
107            nodes.insert(source);
108            nodes.insert(target);
109        }
110
111        Ok(Box::new(nodes.into_iter().map(|node| (node, Ok(())))))
112    }
113}
114
115impl<I: IndexDomain> LaneKernel<Indexed<I, EntityReference<EdgeIndex>>, Single> for NodesOperation {
116    type Output = NodesOperand<Unordered>;
117
118    fn execute<'a>(
119        graphrecord: &'a GraphRecord,
120        value: KeyedStream<'a, I, EntityReference<EdgeIndex>, Single>,
121        _prepared: Self::Prepared<'a>,
122    ) -> QueryResult<<Self::Output as EvaluateOperand>::ReturnValue<'a>> {
123        let Some(value) = value else {
124            return Ok(Box::new(empty()));
125        };
126        let edge = value.1?;
127        let (source, target) = graphrecord.edge_endpoints(edge).expect("Edge must exist");
128        let nodes: GrHashSet<_> = once(source).chain(once(target)).collect();
129
130        Ok(Box::new(nodes.into_iter().map(|node| (node, Ok(())))))
131    }
132}
133
134impl<I: IndexDomain> LaneKernel<Indexed<I, EntityReference<EdgeIndex>>, Definite>
135    for NodesOperation
136{
137    type Output = NodesOperand<Unordered>;
138
139    fn execute<'a>(
140        graphrecord: &'a GraphRecord,
141        value: KeyedStream<'a, I, EntityReference<EdgeIndex>, Definite>,
142        _prepared: Self::Prepared<'a>,
143    ) -> QueryResult<<Self::Output as EvaluateOperand>::ReturnValue<'a>> {
144        let edge = value.1?;
145        let (source, target) = graphrecord.edge_endpoints(edge).expect("Edge must exist");
146        let nodes: GrHashSet<_> = once(source).chain(once(target)).collect();
147
148        Ok(Box::new(nodes.into_iter().map(|node| (node, Ok(())))))
149    }
150}
151
152impl<O: Apply<NodesOperation>> Nodes for O {
153    type ReturnOperand = O::Output;
154
155    fn nodes(&self) -> Self::ReturnOperand {
156        Self::ReturnOperand::new(OperationContext::new(self.clone(), NodesOperation))
157    }
158}
159
160operation_manifest! {
161    NodesOperation {
162        method: Nodes::nodes;
163        scope: lane;
164
165        kernel {
166            parameters: <O: OrderState>;
167            input: (Indexed<EdgeIndex, Unit>, Multiple<O>);
168            output: NodesOperand<Unordered>;
169        }
170        kernel {
171            parameters: <>;
172            input: (Indexed<EdgeIndex, Unit>, Single);
173            output: NodesOperand<Unordered>;
174        }
175        kernel {
176            parameters: <>;
177            input: (Indexed<EdgeIndex, Unit>, Definite);
178            output: NodesOperand<Unordered>;
179        }
180        kernel {
181            parameters: <I: IndexDomain, O: OrderState>;
182            input: (Indexed<I, EntityReference<EdgeIndex>>, Multiple<O>);
183            output: NodesOperand<Unordered>;
184        }
185        kernel {
186            parameters: <I: IndexDomain>;
187            input: (Indexed<I, EntityReference<EdgeIndex>>, Single);
188            output: NodesOperand<Unordered>;
189        }
190        kernel {
191            parameters: <I: IndexDomain>;
192            input: (Indexed<I, EntityReference<EdgeIndex>>, Definite);
193            output: NodesOperand<Unordered>;
194        }
195    }
196}