1use std::collections::BTreeMap;
9use std::sync::Arc;
10
11use lora_compiler::physical::{
12 CallSubqueryExec, ExpandExec, FilterExec, HashAggregationExec, LimitExec, NodeByIdSeekExec,
13 NodeByLabelScanExec, NodeByPointScanExec, NodeByPropertyRangeScanExec, NodeByPropertyScanExec,
14 NodeByTextScanExec, NodeScanExec, OptionalMatchExec, PathBuildExec, PhysicalNodeId, PhysicalOp,
15 PhysicalPlan, ProjectionExec, RelByIdSeekExec, RelByPointScanExec, RelByPropertyRangeScanExec,
16 RelByTextScanExec, SortExec, UnwindExec,
17};
18use lora_compiler::CompiledQuery;
19use lora_store::GraphStorage;
20
21use crate::errors::{ExecResult, ExecutorError};
22use crate::eval::clear_eval_error;
23use crate::executor::{plan_may_need_hydration, ExecutionContext, Executor};
24use crate::profile::wrap_metered;
25use crate::value::{LoraValue, Row};
26
27use super::aggregate::HashAggregationSource;
28use super::call_subquery::CallSubquerySource;
29use super::expand::{ExpandSource, VariableLengthExpandSource};
30use super::filter::FilterSource;
31use super::optional::OptionalMatchSource;
32use super::path::PathBuildSource;
33use super::projection::{DistinctSource, ProjectionSource, UnwindSource};
34use super::scan::{
35 BufferedIndexScanSource, NodeByLabelScanSource, NodeByPropertyScanSource, NodeScanSource,
36};
37use super::sort::{LimitSource, SortSource};
38use super::union::UnionSource;
39use super::{drain, ArgumentSource, BufferedRowSource, HydratingSource, RowSource, StreamCtx};
40
41pub(super) fn compiled_to_streaming<'a, S: GraphStorage + 'a>(
57 compiled: &'a CompiledQuery,
58 storage: &'a S,
59 params: BTreeMap<String, LoraValue>,
60) -> ExecResult<Box<dyn RowSource + 'a>> {
61 let params = Arc::new(params);
62 let _stable = super::scan::StableStoreScope::enter();
65
66 if compiled.unions.is_empty() {
67 let plan = &compiled.physical;
68 let inner = build_streaming(plan, plan.root, storage, params)?;
69 if !plan_may_need_hydration(plan) {
70 return Ok(inner);
71 }
72 return Ok(Box::new(HydratingSource::new(inner, storage)));
73 }
74
75 let mut branches: Vec<Box<dyn RowSource + 'a>> = Vec::with_capacity(compiled.unions.len() + 1);
76
77 let head_inner = build_streaming(
78 &compiled.physical,
79 compiled.physical.root,
80 storage,
81 params.clone(),
82 )?;
83 if plan_may_need_hydration(&compiled.physical) {
84 branches.push(Box::new(HydratingSource::new(head_inner, storage)));
85 } else {
86 branches.push(head_inner);
87 }
88
89 let mut needs_dedup = false;
90 for branch in &compiled.unions {
91 let inner = build_streaming(
92 &branch.physical,
93 branch.physical.root,
94 storage,
95 params.clone(),
96 )?;
97 if plan_may_need_hydration(&branch.physical) {
98 branches.push(Box::new(HydratingSource::new(inner, storage)));
99 } else {
100 branches.push(inner);
101 }
102 if !branch.all {
103 needs_dedup = true;
104 }
105 }
106
107 Ok(Box::new(UnionSource::new(branches, needs_dedup)))
108}
109
110pub(super) fn is_streaming_op(op: &PhysicalOp) -> bool {
118 match op {
119 PhysicalOp::Argument(_)
120 | PhysicalOp::NodeScan(_)
121 | PhysicalOp::NodeByLabelScan(_)
122 | PhysicalOp::NodeByPropertyScan(_)
123 | PhysicalOp::NodeByIdSeek(_)
124 | PhysicalOp::RelByIdSeek(_)
125 | PhysicalOp::NodeByPropertyRangeScan(_)
126 | PhysicalOp::NodeByTextScan(_)
127 | PhysicalOp::NodeByPointScan(_)
128 | PhysicalOp::RelByPropertyRangeScan(_)
129 | PhysicalOp::RelByTextScan(_)
130 | PhysicalOp::RelByPointScan(_)
131 | PhysicalOp::Filter(_)
132 | PhysicalOp::Unwind(_)
133 | PhysicalOp::Limit(_)
134 | PhysicalOp::Sort(_)
141 | PhysicalOp::HashAggregation(_)
142 | PhysicalOp::OptionalMatch(_)
143 | PhysicalOp::PathBuild(_)
144 | PhysicalOp::CallSubquery(_)
145 | PhysicalOp::Projection(_) => true,
149 PhysicalOp::Expand(_) => true,
152 _ => false,
153 }
154}
155
156pub(crate) fn write_op_input(
161 plan: &PhysicalPlan,
162 node_id: PhysicalNodeId,
163) -> Option<PhysicalNodeId> {
164 match &plan.nodes[node_id] {
165 PhysicalOp::Create(o) => Some(o.input),
166 PhysicalOp::Set(o) => Some(o.input),
167 PhysicalOp::Delete(o) => Some(o.input),
168 PhysicalOp::Remove(o) => Some(o.input),
169 PhysicalOp::Merge(o) => Some(o.input),
170 _ => None,
171 }
172}
173
174pub(crate) fn subtree_is_fully_streaming(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
181 let op = &plan.nodes[node_id];
182 if !is_streaming_op(op) {
183 return false;
184 }
185 let child = match op {
186 PhysicalOp::Argument(_) => return true,
187 PhysicalOp::NodeScan(o) => o.input,
188 PhysicalOp::NodeByLabelScan(o) => o.input,
189 PhysicalOp::NodeByPropertyScan(o) => o.input,
190 PhysicalOp::NodeByIdSeek(o) => o.input,
191 PhysicalOp::RelByIdSeek(o) => o.input,
192 PhysicalOp::NodeByPropertyRangeScan(o) => o.input,
193 PhysicalOp::NodeByTextScan(o) => o.input,
194 PhysicalOp::NodeByPointScan(o) => o.input,
195 PhysicalOp::RelByPropertyRangeScan(o) => o.input,
196 PhysicalOp::RelByTextScan(o) => o.input,
197 PhysicalOp::RelByPointScan(o) => o.input,
198 PhysicalOp::Filter(o) => Some(o.input),
199 PhysicalOp::Unwind(o) => Some(o.input),
200 PhysicalOp::Limit(o) => Some(o.input),
201 PhysicalOp::Expand(o) => Some(o.input),
202 PhysicalOp::Projection(o) => Some(o.input),
203 PhysicalOp::Sort(o) => Some(o.input),
204 PhysicalOp::HashAggregation(o) => Some(o.input),
205 PhysicalOp::OptionalMatch(o) => Some(o.input),
206 PhysicalOp::CallSubquery(o) if subtree_has_write(plan, o.inner) => return false,
209 PhysicalOp::CallSubquery(o) => Some(o.input),
210 PhysicalOp::PathBuild(o) => Some(o.input),
211 _ => return false,
213 };
214 match child {
215 None => true,
216 Some(c) => subtree_is_fully_streaming(plan, c),
217 }
218}
219
220pub(crate) fn subtree_has_write(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
223 let op = &plan.nodes[node_id];
224 let (first, second) = match op {
225 PhysicalOp::Create(_)
226 | PhysicalOp::Merge(_)
227 | PhysicalOp::Delete(_)
228 | PhysicalOp::Set(_)
229 | PhysicalOp::Remove(_)
230 | PhysicalOp::Foreach(_) => return true,
231 PhysicalOp::Argument(_) => (None, None),
232 PhysicalOp::NodeScan(o) => (o.input, None),
233 PhysicalOp::NodeByLabelScan(o) => (o.input, None),
234 PhysicalOp::NodeByPropertyScan(o) => (o.input, None),
235 PhysicalOp::NodeByIdSeek(o) => (o.input, None),
236 PhysicalOp::RelByIdSeek(o) => (o.input, None),
237 PhysicalOp::NodeByPropertyRangeScan(o) => (o.input, None),
238 PhysicalOp::NodeByTextScan(o) => (o.input, None),
239 PhysicalOp::NodeByPointScan(o) => (o.input, None),
240 PhysicalOp::RelByPropertyRangeScan(o) => (o.input, None),
241 PhysicalOp::RelByTextScan(o) => (o.input, None),
242 PhysicalOp::RelByPointScan(o) => (o.input, None),
243 PhysicalOp::Expand(o) => (Some(o.input), None),
244 PhysicalOp::Filter(o) => (Some(o.input), None),
245 PhysicalOp::Projection(o) => (Some(o.input), None),
246 PhysicalOp::Unwind(o) => (Some(o.input), None),
247 PhysicalOp::HashAggregation(o) => (Some(o.input), None),
248 PhysicalOp::Sort(o) => (Some(o.input), None),
249 PhysicalOp::Limit(o) => (Some(o.input), None),
250 PhysicalOp::PathBuild(o) => (Some(o.input), None),
251 PhysicalOp::OptionalMatch(o) => (Some(o.input), Some(o.inner)),
252 PhysicalOp::CallSubquery(o) => (Some(o.input), Some(o.inner)),
253 };
254 first.is_some_and(|c| subtree_has_write(plan, c))
255 || second.is_some_and(|c| subtree_has_write(plan, c))
256}
257
258pub(crate) fn build_streaming<'a, S: GraphStorage + 'a>(
259 plan: &'a PhysicalPlan,
260 node_id: PhysicalNodeId,
261 storage: &'a S,
262 params: Arc<BTreeMap<String, LoraValue>>,
263) -> ExecResult<Box<dyn RowSource + 'a>> {
264 build_streaming_dispatch(plan, node_id, storage, params, None)
265}
266
267pub(crate) fn build_streaming_seeded<'a, S: GraphStorage + 'a>(
271 plan: &'a PhysicalPlan,
272 node_id: PhysicalNodeId,
273 storage: &'a S,
274 params: Arc<BTreeMap<String, LoraValue>>,
275 seed: Row,
276) -> ExecResult<Box<dyn RowSource + 'a>> {
277 build_streaming_dispatch(plan, node_id, storage, params, Some(seed))
278}
279
280fn build_streaming_dispatch<'a, S: GraphStorage + 'a>(
281 plan: &'a PhysicalPlan,
282 node_id: PhysicalNodeId,
283 storage: &'a S,
284 params: Arc<BTreeMap<String, LoraValue>>,
285 seed: Option<Row>,
286) -> ExecResult<Box<dyn RowSource + 'a>> {
287 build_streaming_inner(plan, node_id, storage, params, seed).map(|src| {
288 super::source::DeadlineSource::wrap(
289 wrap_metered(node_id, src),
290 crate::cancel::active_deadline(),
291 )
292 })
293}
294
295fn build_streaming_inner<'a, S: GraphStorage + 'a>(
296 plan: &'a PhysicalPlan,
297 node_id: PhysicalNodeId,
298 storage: &'a S,
299 params: Arc<BTreeMap<String, LoraValue>>,
300 seed: Option<Row>,
301) -> ExecResult<Box<dyn RowSource + 'a>> {
302 let op = &plan.nodes[node_id];
303
304 if !is_streaming_op(op) {
305 return build_buffered_subtree(plan, node_id, storage, ¶ms);
306 }
307
308 match op {
309 PhysicalOp::Argument(_) => match seed {
310 Some(seed_row) => Ok(Box::new(BufferedRowSource::new(vec![seed_row]))),
311 None => Ok(Box::new(ArgumentSource::new())),
312 },
313
314 PhysicalOp::NodeScan(NodeScanExec { input, var }) => {
315 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
316 Ok(Box::new(NodeScanSource::new(upstream, storage, *var)))
317 }
318
319 PhysicalOp::NodeByLabelScan(NodeByLabelScanExec { input, var, labels }) => {
320 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
321 Ok(Box::new(NodeByLabelScanSource::new(
322 upstream, storage, *var, labels,
323 )))
324 }
325
326 PhysicalOp::NodeByPropertyScan(NodeByPropertyScanExec {
327 input,
328 var,
329 labels,
330 key,
331 value,
332 in_list,
333 }) => {
334 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
335 let ctx = StreamCtx::new(storage, params);
336 Ok(Box::new(NodeByPropertyScanSource::new(
337 upstream, ctx, *var, labels, key, value, *in_list,
338 )))
339 }
340
341 PhysicalOp::NodeByIdSeek(op @ NodeByIdSeekExec { input, .. }) => {
342 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
343 let ctx = StreamCtx::new(storage, params);
344 Ok(Box::new(BufferedIndexScanSource::node_id(
345 upstream, ctx, op,
346 )))
347 }
348
349 PhysicalOp::RelByIdSeek(op @ RelByIdSeekExec { input, .. }) => {
350 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
351 let ctx = StreamCtx::new(storage, params);
352 Ok(Box::new(BufferedIndexScanSource::rel_id(upstream, ctx, op)))
353 }
354
355 PhysicalOp::NodeByPropertyRangeScan(op @ NodeByPropertyRangeScanExec { input, .. }) => {
356 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
357 let ctx = StreamCtx::new(storage, params);
358 if op.order.is_some() {
359 return Ok(Box::new(super::scan::OrderedRangeScanSource::new(
360 upstream, ctx, op,
361 )));
362 }
363 Ok(Box::new(BufferedIndexScanSource::node_range(
364 upstream, ctx, op,
365 )))
366 }
367
368 PhysicalOp::NodeByTextScan(op @ NodeByTextScanExec { input, .. }) => {
369 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
370 let ctx = StreamCtx::new(storage, params);
371 Ok(Box::new(BufferedIndexScanSource::node_text(
372 upstream, ctx, op,
373 )))
374 }
375
376 PhysicalOp::NodeByPointScan(op @ NodeByPointScanExec { input, .. }) => {
377 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
378 let ctx = StreamCtx::new(storage, params);
379 Ok(Box::new(BufferedIndexScanSource::node_point(
380 upstream, ctx, op,
381 )))
382 }
383
384 PhysicalOp::RelByPropertyRangeScan(op @ RelByPropertyRangeScanExec { input, .. }) => {
385 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
386 let ctx = StreamCtx::new(storage, params);
387 Ok(Box::new(BufferedIndexScanSource::rel_range(
388 upstream, ctx, op,
389 )))
390 }
391
392 PhysicalOp::RelByTextScan(op @ RelByTextScanExec { input, .. }) => {
393 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
394 let ctx = StreamCtx::new(storage, params);
395 Ok(Box::new(BufferedIndexScanSource::rel_text(
396 upstream, ctx, op,
397 )))
398 }
399
400 PhysicalOp::RelByPointScan(op @ RelByPointScanExec { input, .. }) => {
401 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
402 let ctx = StreamCtx::new(storage, params);
403 Ok(Box::new(BufferedIndexScanSource::rel_point(
404 upstream, ctx, op,
405 )))
406 }
407
408 PhysicalOp::Expand(ExpandExec {
409 input,
410 src,
411 rel,
412 dst,
413 types,
414 direction,
415 rel_properties,
416 range,
417 }) => {
418 let upstream =
419 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
420 let ctx = StreamCtx::new(storage, params);
421 match range.as_ref() {
422 Some(range) => Ok(Box::new(VariableLengthExpandSource::new(
423 upstream, ctx, *src, *rel, *dst, types, *direction, range,
424 ))),
425 None => Ok(Box::new(ExpandSource::new(
426 upstream,
427 ctx,
428 *src,
429 *rel,
430 *dst,
431 types,
432 *direction,
433 rel_properties.as_ref(),
434 ))),
435 }
436 }
437
438 PhysicalOp::Filter(FilterExec { input, predicate }) => {
439 let upstream =
440 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
441 let ctx = StreamCtx::new(storage, params);
442 Ok(Box::new(FilterSource::new(upstream, ctx, predicate)))
443 }
444
445 PhysicalOp::Projection(ProjectionExec {
446 input,
447 distinct,
448 items,
449 include_existing,
450 }) => {
451 let upstream =
452 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
453 let ctx = StreamCtx::new(storage, params);
454 let proj: Box<dyn RowSource + 'a> = Box::new(ProjectionSource::new(
455 upstream,
456 ctx,
457 items,
458 *include_existing,
459 ));
460 if *distinct {
461 Ok(Box::new(DistinctSource::new(proj)))
462 } else {
463 Ok(proj)
464 }
465 }
466
467 PhysicalOp::Unwind(UnwindExec { input, expr, alias }) => {
468 let upstream =
469 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
470 let ctx = StreamCtx::new(storage, params);
471 Ok(Box::new(UnwindSource::new(upstream, ctx, expr, *alias)))
472 }
473
474 PhysicalOp::Limit(LimitExec { input, skip, limit }) => {
475 let upstream =
476 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
477 let ctx = StreamCtx::new(storage, params);
480 let eval_ctx = ctx.eval_ctx();
481 let skip_n = match skip.as_ref() {
482 Some(e) => crate::executor::eval_row_count("SKIP", e, &eval_ctx)?,
483 None => 0,
484 };
485 let limit_n = match limit.as_ref() {
486 Some(e) => Some(crate::executor::eval_row_count("LIMIT", e, &eval_ctx)?),
487 None => None,
488 };
489 Ok(Box::new(LimitSource::new(upstream, skip_n, limit_n)))
490 }
491
492 PhysicalOp::Sort(SortExec {
493 input,
494 items,
495 top_k,
496 limit,
497 }) => {
498 let upstream =
499 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
500 let ctx = StreamCtx::new(storage, params);
501 let bound = crate::executor::sort_row_bound(*top_k, limit.as_ref(), &ctx.eval_ctx());
502 Ok(Box::new(SortSource::new_with_top_k(
503 upstream, ctx, items, bound,
504 )))
505 }
506
507 PhysicalOp::HashAggregation(
508 agg @ HashAggregationExec {
509 input,
510 group_by,
511 aggregates,
512 },
513 ) => {
514 if seed.is_none() {
518 if let Some(rows) =
519 crate::executor::count_all_scan_aggregation_rows(storage, plan, agg)
520 {
521 return Ok(Box::new(BufferedRowSource::new(rows)));
522 }
523 }
524 let upstream =
525 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
526 let ctx = StreamCtx::new(storage, params);
527 Ok(Box::new(HashAggregationSource::new(
528 upstream, ctx, group_by, aggregates,
529 )))
530 }
531
532 PhysicalOp::OptionalMatch(OptionalMatchExec {
533 input,
534 inner,
535 new_vars,
536 }) => {
537 let upstream =
538 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
539 let ctx = StreamCtx::new(storage, params);
540 Ok(Box::new(OptionalMatchSource::new(
541 upstream, ctx, plan, *inner, new_vars,
542 )))
543 }
544
545 PhysicalOp::CallSubquery(CallSubqueryExec {
546 input,
547 inner,
548 new_vars,
549 }) => {
550 let upstream = build_streaming_dispatch(plan, *input, storage, params.clone(), seed)?;
551 Ok(Box::new(CallSubquerySource::new(
552 upstream, plan, *inner, storage, params, new_vars,
553 )))
554 }
555
556 PhysicalOp::PathBuild(PathBuildExec {
557 input,
558 output,
559 node_vars,
560 rel_vars,
561 shortest_path_all,
562 }) => {
563 let upstream =
564 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
565 let ctx = StreamCtx::new(storage, params);
566 Ok(Box::new(PathBuildSource::new(
567 upstream,
568 ctx,
569 *output,
570 node_vars,
571 rel_vars,
572 *shortest_path_all,
573 )))
574 }
575
576 _ => Err(ExecutorError::RuntimeError(format!(
579 "non-streaming op reached streaming branch: {op:?}"
580 ))),
581 }
582}
583
584fn open_input<'a, S: GraphStorage + 'a>(
588 plan: &'a PhysicalPlan,
589 input: Option<PhysicalNodeId>,
590 storage: &'a S,
591 params: Arc<BTreeMap<String, LoraValue>>,
592 seed: Option<Row>,
593) -> ExecResult<Box<dyn RowSource + 'a>> {
594 match input {
595 Some(input) => build_streaming_dispatch(plan, input, storage, params, seed),
596 None => match seed {
597 Some(seed_row) => Ok(Box::new(BufferedRowSource::new(vec![seed_row]))),
598 None => Ok(Box::new(ArgumentSource::new())),
599 },
600 }
601}
602
603fn build_buffered_subtree<'a, S: GraphStorage + 'a>(
610 plan: &'a PhysicalPlan,
611 node_id: PhysicalNodeId,
612 storage: &'a S,
613 params: &Arc<BTreeMap<String, LoraValue>>,
614) -> ExecResult<Box<dyn RowSource + 'a>> {
615 let executor = Executor::with_deadline(
619 ExecutionContext {
620 storage,
621 params: (**params).clone(),
622 },
623 crate::cancel::active_deadline(),
624 );
625 let rows = executor.execute_subtree(plan, node_id)?;
626 Ok(Box::new(BufferedRowSource::new(rows)))
627}
628
629pub struct PullExecutor<'a, S: GraphStorage> {
635 storage: &'a S,
636 params: BTreeMap<String, LoraValue>,
637}
638
639impl<'a, S: GraphStorage> PullExecutor<'a, S> {
640 pub fn new(storage: &'a S, params: BTreeMap<String, LoraValue>) -> Self {
641 Self { storage, params }
642 }
643
644 pub fn open_compiled(self, compiled: &'a CompiledQuery) -> ExecResult<Box<dyn RowSource + 'a>>
653 where
654 S: 'a,
655 {
656 clear_eval_error();
657 compiled_to_streaming(compiled, self.storage, self.params)
658 }
659}
660
661pub fn collect_compiled<'a, S: GraphStorage + 'a>(
664 storage: &'a S,
665 params: BTreeMap<String, LoraValue>,
666 compiled: &'a CompiledQuery,
667) -> ExecResult<Vec<Row>> {
668 let mut cursor = PullExecutor::new(storage, params).open_compiled(compiled)?;
669 drain(cursor.as_mut())
670}
671
672pub fn collect_compiled_with_deadline<'a, S: GraphStorage + 'a>(
679 storage: &'a S,
680 params: BTreeMap<String, LoraValue>,
681 compiled: &'a CompiledQuery,
682 deadline: Option<web_time::Instant>,
683) -> ExecResult<Vec<Row>> {
684 let _deadline_scope = crate::cancel::DeadlineScope::enter(deadline);
685 if let Some(deadline) = deadline {
686 if crate::cancel::deadline_reached(deadline) {
687 return Err(ExecutorError::QueryTimeout);
688 }
689 }
690 let mut cursor = PullExecutor::new(storage, params).open_compiled(compiled)?;
691 drain(cursor.as_mut())
692}