1use std::collections::BTreeMap;
9use std::sync::Arc;
10
11use lora_compiler::physical::{
12 ExpandExec, FilterExec, HashAggregationExec, LimitExec, NodeByLabelScanExec,
13 NodeByPropertyScanExec, NodeScanExec, OptionalMatchExec, PathBuildExec, PhysicalNodeId,
14 PhysicalOp, PhysicalPlan, ProjectionExec, SortExec, UnwindExec,
15};
16use lora_compiler::CompiledQuery;
17use lora_store::GraphStorage;
18
19use crate::errors::ExecResult;
20use crate::eval::{clear_eval_error, eval_expr};
21use crate::executor::{ExecutionContext, Executor};
22use crate::profile::wrap_metered;
23use crate::value::{LoraValue, Row};
24
25use super::aggregate::HashAggregationSource;
26use super::expand::{ExpandSource, VariableLengthExpandSource};
27use super::filter::FilterSource;
28use super::optional::OptionalMatchSource;
29use super::path::PathBuildSource;
30use super::projection::{DistinctSource, ProjectionSource, UnwindSource};
31use super::scan::{NodeByLabelScanSource, NodeByPropertyScanSource, NodeScanSource};
32use super::sort::{LimitSource, SortSource};
33use super::union::UnionSource;
34use super::{drain, ArgumentSource, BufferedRowSource, HydratingSource, RowSource, StreamCtx};
35
36pub(super) fn compiled_to_streaming<'a, S: GraphStorage + 'a>(
52 compiled: &'a CompiledQuery,
53 storage: &'a S,
54 params: BTreeMap<String, LoraValue>,
55) -> ExecResult<Box<dyn RowSource + 'a>> {
56 let params = Arc::new(params);
57
58 if compiled.unions.is_empty() {
59 let plan = &compiled.physical;
60 let inner = build_streaming(plan, plan.root, storage, params)?;
61 return Ok(Box::new(HydratingSource::new(inner, storage)));
62 }
63
64 let mut branches: Vec<Box<dyn RowSource + 'a>> = Vec::with_capacity(compiled.unions.len() + 1);
65
66 let head_inner = build_streaming(
67 &compiled.physical,
68 compiled.physical.root,
69 storage,
70 params.clone(),
71 )?;
72 branches.push(Box::new(HydratingSource::new(head_inner, storage)));
73
74 let mut needs_dedup = false;
75 for branch in &compiled.unions {
76 let inner = build_streaming(
77 &branch.physical,
78 branch.physical.root,
79 storage,
80 params.clone(),
81 )?;
82 branches.push(Box::new(HydratingSource::new(inner, storage)));
83 if !branch.all {
84 needs_dedup = true;
85 }
86 }
87
88 Ok(Box::new(UnionSource::new(branches, needs_dedup)))
89}
90
91pub(super) fn is_streaming_op(op: &PhysicalOp) -> bool {
99 match op {
100 PhysicalOp::Argument(_)
101 | PhysicalOp::NodeScan(_)
102 | PhysicalOp::NodeByLabelScan(_)
103 | PhysicalOp::NodeByPropertyScan(_)
104 | PhysicalOp::Filter(_)
105 | PhysicalOp::Unwind(_)
106 | PhysicalOp::Limit(_)
107 | PhysicalOp::Sort(_)
114 | PhysicalOp::HashAggregation(_)
115 | PhysicalOp::OptionalMatch(_)
116 | PhysicalOp::PathBuild(_)
117 | PhysicalOp::Projection(_) => true,
121 PhysicalOp::Expand(_) => true,
124 _ => false,
125 }
126}
127
128pub(crate) fn write_op_input(
133 plan: &PhysicalPlan,
134 node_id: PhysicalNodeId,
135) -> Option<PhysicalNodeId> {
136 match &plan.nodes[node_id] {
137 PhysicalOp::Create(o) => Some(o.input),
138 PhysicalOp::Set(o) => Some(o.input),
139 PhysicalOp::Delete(o) => Some(o.input),
140 PhysicalOp::Remove(o) => Some(o.input),
141 PhysicalOp::Merge(o) => Some(o.input),
142 _ => None,
143 }
144}
145
146pub(crate) fn subtree_is_fully_streaming(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
153 let op = &plan.nodes[node_id];
154 if !is_streaming_op(op) {
155 return false;
156 }
157 let child = match op {
158 PhysicalOp::Argument(_) => return true,
159 PhysicalOp::NodeScan(o) => o.input,
160 PhysicalOp::NodeByLabelScan(o) => o.input,
161 PhysicalOp::NodeByPropertyScan(o) => o.input,
162 PhysicalOp::Filter(o) => Some(o.input),
163 PhysicalOp::Unwind(o) => Some(o.input),
164 PhysicalOp::Limit(o) => Some(o.input),
165 PhysicalOp::Expand(o) => Some(o.input),
166 PhysicalOp::Projection(o) => Some(o.input),
167 PhysicalOp::Sort(o) => Some(o.input),
168 PhysicalOp::HashAggregation(o) => Some(o.input),
169 PhysicalOp::OptionalMatch(o) => Some(o.input),
170 PhysicalOp::PathBuild(o) => Some(o.input),
171 _ => return false,
173 };
174 match child {
175 None => true,
176 Some(c) => subtree_is_fully_streaming(plan, c),
177 }
178}
179
180pub(crate) fn build_streaming<'a, S: GraphStorage + 'a>(
181 plan: &'a PhysicalPlan,
182 node_id: PhysicalNodeId,
183 storage: &'a S,
184 params: Arc<BTreeMap<String, LoraValue>>,
185) -> ExecResult<Box<dyn RowSource + 'a>> {
186 build_streaming_inner(plan, node_id, storage, params).map(|src| wrap_metered(node_id, src))
187}
188
189fn build_streaming_inner<'a, S: GraphStorage + 'a>(
190 plan: &'a PhysicalPlan,
191 node_id: PhysicalNodeId,
192 storage: &'a S,
193 params: Arc<BTreeMap<String, LoraValue>>,
194) -> ExecResult<Box<dyn RowSource + 'a>> {
195 let op = &plan.nodes[node_id];
196
197 if !is_streaming_op(op) {
198 return build_buffered_subtree(plan, node_id, storage, ¶ms);
199 }
200
201 match op {
202 PhysicalOp::Argument(_) => Ok(Box::new(ArgumentSource::new())),
203
204 PhysicalOp::NodeScan(NodeScanExec { input, var }) => {
205 let upstream = open_input(plan, *input, storage, params.clone())?;
206 Ok(Box::new(NodeScanSource::new(upstream, storage, *var)))
207 }
208
209 PhysicalOp::NodeByLabelScan(NodeByLabelScanExec { input, var, labels }) => {
210 let upstream = open_input(plan, *input, storage, params.clone())?;
211 Ok(Box::new(NodeByLabelScanSource::new(
212 upstream, storage, *var, labels,
213 )))
214 }
215
216 PhysicalOp::NodeByPropertyScan(NodeByPropertyScanExec {
217 input,
218 var,
219 labels,
220 key,
221 value,
222 }) => {
223 let upstream = open_input(plan, *input, storage, params.clone())?;
224 let ctx = StreamCtx::new(storage, params);
225 Ok(Box::new(NodeByPropertyScanSource::new(
226 upstream, ctx, *var, labels, key, value,
227 )))
228 }
229
230 PhysicalOp::Expand(ExpandExec {
231 input,
232 src,
233 rel,
234 dst,
235 types,
236 direction,
237 rel_properties,
238 range,
239 }) => {
240 let upstream = build_streaming(plan, *input, storage, params.clone())?;
241 let ctx = StreamCtx::new(storage, params);
242 match range.as_ref() {
243 Some(range) => Ok(Box::new(VariableLengthExpandSource::new(
244 upstream, ctx, *src, *rel, *dst, types, *direction, range,
245 ))),
246 None => Ok(Box::new(ExpandSource::new(
247 upstream,
248 ctx,
249 *src,
250 *rel,
251 *dst,
252 types,
253 *direction,
254 rel_properties.as_ref(),
255 ))),
256 }
257 }
258
259 PhysicalOp::Filter(FilterExec { input, predicate }) => {
260 let upstream = build_streaming(plan, *input, storage, params.clone())?;
261 let ctx = StreamCtx::new(storage, params);
262 Ok(Box::new(FilterSource::new(upstream, ctx, predicate)))
263 }
264
265 PhysicalOp::Projection(ProjectionExec {
266 input,
267 distinct,
268 items,
269 include_existing,
270 }) => {
271 let upstream = build_streaming(plan, *input, storage, params.clone())?;
272 let ctx = StreamCtx::new(storage, params);
273 let proj: Box<dyn RowSource + 'a> = Box::new(ProjectionSource::new(
274 upstream,
275 ctx,
276 items,
277 *include_existing,
278 ));
279 if *distinct {
280 Ok(Box::new(DistinctSource::new(proj)))
281 } else {
282 Ok(proj)
283 }
284 }
285
286 PhysicalOp::Unwind(UnwindExec { input, expr, alias }) => {
287 let upstream = build_streaming(plan, *input, storage, params.clone())?;
288 let ctx = StreamCtx::new(storage, params);
289 Ok(Box::new(UnwindSource::new(upstream, ctx, expr, *alias)))
290 }
291
292 PhysicalOp::Limit(LimitExec { input, skip, limit }) => {
293 let upstream = build_streaming(plan, *input, storage, params.clone())?;
294 let ctx = StreamCtx::new(storage, params);
297 let eval_ctx = ctx.eval_ctx();
298 let scratch = Row::new();
299 let skip_n = skip
300 .as_ref()
301 .and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
302 .unwrap_or(0)
303 .max(0) as usize;
304 let limit_n = limit
305 .as_ref()
306 .and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
307 .map(|n| n.max(0) as usize);
308 Ok(Box::new(LimitSource::new(upstream, skip_n, limit_n)))
309 }
310
311 PhysicalOp::Sort(SortExec { input, items }) => {
312 let upstream = build_streaming(plan, *input, storage, params.clone())?;
313 let ctx = StreamCtx::new(storage, params);
314 Ok(Box::new(SortSource::new(upstream, ctx, items)))
315 }
316
317 PhysicalOp::HashAggregation(HashAggregationExec {
318 input,
319 group_by,
320 aggregates,
321 }) => {
322 let upstream = build_streaming(plan, *input, storage, params.clone())?;
323 let ctx = StreamCtx::new(storage, params);
324 Ok(Box::new(HashAggregationSource::new(
325 upstream, ctx, group_by, aggregates,
326 )))
327 }
328
329 PhysicalOp::OptionalMatch(OptionalMatchExec {
330 input,
331 inner,
332 new_vars,
333 }) => {
334 let upstream = build_streaming(plan, *input, storage, params.clone())?;
335 let ctx = StreamCtx::new(storage, params);
336 Ok(Box::new(OptionalMatchSource::new(
337 upstream, ctx, plan, *inner, new_vars,
338 )))
339 }
340
341 PhysicalOp::PathBuild(PathBuildExec {
342 input,
343 output,
344 node_vars,
345 rel_vars,
346 shortest_path_all,
347 }) => {
348 let upstream = build_streaming(plan, *input, storage, params.clone())?;
349 let ctx = StreamCtx::new(storage, params);
350 Ok(Box::new(PathBuildSource::new(
351 upstream,
352 ctx,
353 *output,
354 node_vars,
355 rel_vars,
356 *shortest_path_all,
357 )))
358 }
359
360 _ => unreachable!("non-streaming op reached streaming branch: {op:?}"),
362 }
363}
364
365fn open_input<'a, S: GraphStorage + 'a>(
369 plan: &'a PhysicalPlan,
370 input: Option<PhysicalNodeId>,
371 storage: &'a S,
372 params: Arc<BTreeMap<String, LoraValue>>,
373) -> ExecResult<Box<dyn RowSource + 'a>> {
374 match input {
375 Some(input) => build_streaming(plan, input, storage, params),
376 None => Ok(Box::new(ArgumentSource::new())),
377 }
378}
379
380fn build_buffered_subtree<'a, S: GraphStorage + 'a>(
387 plan: &'a PhysicalPlan,
388 node_id: PhysicalNodeId,
389 storage: &'a S,
390 params: &Arc<BTreeMap<String, LoraValue>>,
391) -> ExecResult<Box<dyn RowSource + 'a>> {
392 let executor = Executor::new(ExecutionContext {
396 storage,
397 params: (**params).clone(),
398 });
399 let rows = executor.execute_subtree(plan, node_id)?;
400 Ok(Box::new(BufferedRowSource::new(rows)))
401}
402
403pub struct PullExecutor<'a, S: GraphStorage> {
409 storage: &'a S,
410 params: BTreeMap<String, LoraValue>,
411}
412
413impl<'a, S: GraphStorage> PullExecutor<'a, S> {
414 pub fn new(storage: &'a S, params: BTreeMap<String, LoraValue>) -> Self {
415 Self { storage, params }
416 }
417
418 pub fn open_compiled(self, compiled: &'a CompiledQuery) -> ExecResult<Box<dyn RowSource + 'a>>
427 where
428 S: 'a,
429 {
430 clear_eval_error();
431 compiled_to_streaming(compiled, self.storage, self.params)
432 }
433}
434
435pub fn collect_compiled<'a, S: GraphStorage + 'a>(
438 storage: &'a S,
439 params: BTreeMap<String, LoraValue>,
440 compiled: &'a CompiledQuery,
441) -> ExecResult<Vec<Row>> {
442 let mut cursor = PullExecutor::new(storage, params).open_compiled(compiled)?;
443 drain(cursor.as_mut())
444}