graphrecords_query/operations/traversal/
nodes.rs1use 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}