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