1use std::collections::BTreeMap;
9use std::sync::Arc;
10
11use lora_compiler::physical::{
12 CallSubqueryExec, ExpandExec, FilterExec, HashAggregationExec, LimitExec, NodeByLabelScanExec,
13 NodeByPointScanExec, NodeByPropertyRangeScanExec, NodeByPropertyScanExec, NodeByTextScanExec,
14 NodeScanExec, OptionalMatchExec, PathBuildExec, PhysicalNodeId, PhysicalOp, PhysicalPlan,
15 ProjectionExec, RelByPointScanExec, RelByPropertyRangeScanExec, RelByTextScanExec, SortExec,
16 UnwindExec,
17};
18use lora_compiler::CompiledQuery;
19use lora_store::GraphStorage;
20
21use crate::errors::{ExecResult, ExecutorError};
22use crate::eval::{clear_eval_error, eval_expr};
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::NodeByPropertyRangeScan(_)
121 | PhysicalOp::NodeByTextScan(_)
122 | PhysicalOp::NodeByPointScan(_)
123 | PhysicalOp::RelByPropertyRangeScan(_)
124 | PhysicalOp::RelByTextScan(_)
125 | PhysicalOp::RelByPointScan(_)
126 | PhysicalOp::Filter(_)
127 | PhysicalOp::Unwind(_)
128 | PhysicalOp::Limit(_)
129 | PhysicalOp::Sort(_)
136 | PhysicalOp::HashAggregation(_)
137 | PhysicalOp::OptionalMatch(_)
138 | PhysicalOp::PathBuild(_)
139 | PhysicalOp::CallSubquery(_)
140 | PhysicalOp::Projection(_) => true,
144 PhysicalOp::Expand(_) => true,
147 _ => false,
148 }
149}
150
151pub(crate) fn write_op_input(
156 plan: &PhysicalPlan,
157 node_id: PhysicalNodeId,
158) -> Option<PhysicalNodeId> {
159 match &plan.nodes[node_id] {
160 PhysicalOp::Create(o) => Some(o.input),
161 PhysicalOp::Set(o) => Some(o.input),
162 PhysicalOp::Delete(o) => Some(o.input),
163 PhysicalOp::Remove(o) => Some(o.input),
164 PhysicalOp::Merge(o) => Some(o.input),
165 _ => None,
166 }
167}
168
169pub(crate) fn subtree_is_fully_streaming(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
176 let op = &plan.nodes[node_id];
177 if !is_streaming_op(op) {
178 return false;
179 }
180 let child = match op {
181 PhysicalOp::Argument(_) => return true,
182 PhysicalOp::NodeScan(o) => o.input,
183 PhysicalOp::NodeByLabelScan(o) => o.input,
184 PhysicalOp::NodeByPropertyScan(o) => o.input,
185 PhysicalOp::NodeByPropertyRangeScan(o) => o.input,
186 PhysicalOp::NodeByTextScan(o) => o.input,
187 PhysicalOp::NodeByPointScan(o) => o.input,
188 PhysicalOp::RelByPropertyRangeScan(o) => o.input,
189 PhysicalOp::RelByTextScan(o) => o.input,
190 PhysicalOp::RelByPointScan(o) => o.input,
191 PhysicalOp::Filter(o) => Some(o.input),
192 PhysicalOp::Unwind(o) => Some(o.input),
193 PhysicalOp::Limit(o) => Some(o.input),
194 PhysicalOp::Expand(o) => Some(o.input),
195 PhysicalOp::Projection(o) => Some(o.input),
196 PhysicalOp::Sort(o) => Some(o.input),
197 PhysicalOp::HashAggregation(o) => Some(o.input),
198 PhysicalOp::OptionalMatch(o) => Some(o.input),
199 PhysicalOp::CallSubquery(o) if subtree_has_write(plan, o.inner) => return false,
202 PhysicalOp::CallSubquery(o) => Some(o.input),
203 PhysicalOp::PathBuild(o) => Some(o.input),
204 _ => return false,
206 };
207 match child {
208 None => true,
209 Some(c) => subtree_is_fully_streaming(plan, c),
210 }
211}
212
213pub(crate) fn subtree_has_write(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
216 let op = &plan.nodes[node_id];
217 let (first, second) = match op {
218 PhysicalOp::Create(_)
219 | PhysicalOp::Merge(_)
220 | PhysicalOp::Delete(_)
221 | PhysicalOp::Set(_)
222 | PhysicalOp::Remove(_)
223 | PhysicalOp::Foreach(_) => return true,
224 PhysicalOp::Argument(_) => (None, None),
225 PhysicalOp::NodeScan(o) => (o.input, None),
226 PhysicalOp::NodeByLabelScan(o) => (o.input, None),
227 PhysicalOp::NodeByPropertyScan(o) => (o.input, None),
228 PhysicalOp::NodeByPropertyRangeScan(o) => (o.input, None),
229 PhysicalOp::NodeByTextScan(o) => (o.input, None),
230 PhysicalOp::NodeByPointScan(o) => (o.input, None),
231 PhysicalOp::RelByPropertyRangeScan(o) => (o.input, None),
232 PhysicalOp::RelByTextScan(o) => (o.input, None),
233 PhysicalOp::RelByPointScan(o) => (o.input, None),
234 PhysicalOp::Expand(o) => (Some(o.input), None),
235 PhysicalOp::Filter(o) => (Some(o.input), None),
236 PhysicalOp::Projection(o) => (Some(o.input), None),
237 PhysicalOp::Unwind(o) => (Some(o.input), None),
238 PhysicalOp::HashAggregation(o) => (Some(o.input), None),
239 PhysicalOp::Sort(o) => (Some(o.input), None),
240 PhysicalOp::Limit(o) => (Some(o.input), None),
241 PhysicalOp::PathBuild(o) => (Some(o.input), None),
242 PhysicalOp::OptionalMatch(o) => (Some(o.input), Some(o.inner)),
243 PhysicalOp::CallSubquery(o) => (Some(o.input), Some(o.inner)),
244 };
245 first.is_some_and(|c| subtree_has_write(plan, c))
246 || second.is_some_and(|c| subtree_has_write(plan, c))
247}
248
249pub(crate) fn build_streaming<'a, S: GraphStorage + 'a>(
250 plan: &'a PhysicalPlan,
251 node_id: PhysicalNodeId,
252 storage: &'a S,
253 params: Arc<BTreeMap<String, LoraValue>>,
254) -> ExecResult<Box<dyn RowSource + 'a>> {
255 build_streaming_dispatch(plan, node_id, storage, params, None)
256}
257
258pub(crate) fn build_streaming_seeded<'a, S: GraphStorage + 'a>(
262 plan: &'a PhysicalPlan,
263 node_id: PhysicalNodeId,
264 storage: &'a S,
265 params: Arc<BTreeMap<String, LoraValue>>,
266 seed: Row,
267) -> ExecResult<Box<dyn RowSource + 'a>> {
268 build_streaming_dispatch(plan, node_id, storage, params, Some(seed))
269}
270
271fn build_streaming_dispatch<'a, S: GraphStorage + 'a>(
272 plan: &'a PhysicalPlan,
273 node_id: PhysicalNodeId,
274 storage: &'a S,
275 params: Arc<BTreeMap<String, LoraValue>>,
276 seed: Option<Row>,
277) -> ExecResult<Box<dyn RowSource + 'a>> {
278 build_streaming_inner(plan, node_id, storage, params, seed).map(|src| {
279 super::source::DeadlineSource::wrap(
280 wrap_metered(node_id, src),
281 crate::cancel::active_deadline(),
282 )
283 })
284}
285
286fn build_streaming_inner<'a, S: GraphStorage + 'a>(
287 plan: &'a PhysicalPlan,
288 node_id: PhysicalNodeId,
289 storage: &'a S,
290 params: Arc<BTreeMap<String, LoraValue>>,
291 seed: Option<Row>,
292) -> ExecResult<Box<dyn RowSource + 'a>> {
293 let op = &plan.nodes[node_id];
294
295 if !is_streaming_op(op) {
296 return build_buffered_subtree(plan, node_id, storage, ¶ms);
297 }
298
299 match op {
300 PhysicalOp::Argument(_) => match seed {
301 Some(seed_row) => Ok(Box::new(BufferedRowSource::new(vec![seed_row]))),
302 None => Ok(Box::new(ArgumentSource::new())),
303 },
304
305 PhysicalOp::NodeScan(NodeScanExec { input, var }) => {
306 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
307 Ok(Box::new(NodeScanSource::new(upstream, storage, *var)))
308 }
309
310 PhysicalOp::NodeByLabelScan(NodeByLabelScanExec { input, var, labels }) => {
311 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
312 Ok(Box::new(NodeByLabelScanSource::new(
313 upstream, storage, *var, labels,
314 )))
315 }
316
317 PhysicalOp::NodeByPropertyScan(NodeByPropertyScanExec {
318 input,
319 var,
320 labels,
321 key,
322 value,
323 in_list,
324 }) => {
325 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
326 let ctx = StreamCtx::new(storage, params);
327 Ok(Box::new(NodeByPropertyScanSource::new(
328 upstream, ctx, *var, labels, key, value, *in_list,
329 )))
330 }
331
332 PhysicalOp::NodeByPropertyRangeScan(op @ NodeByPropertyRangeScanExec { input, .. }) => {
333 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
334 let ctx = StreamCtx::new(storage, params);
335 if op.order.is_some() {
336 return Ok(Box::new(super::scan::OrderedRangeScanSource::new(
337 upstream, ctx, op,
338 )));
339 }
340 Ok(Box::new(BufferedIndexScanSource::node_range(
341 upstream, ctx, op,
342 )))
343 }
344
345 PhysicalOp::NodeByTextScan(op @ NodeByTextScanExec { input, .. }) => {
346 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
347 let ctx = StreamCtx::new(storage, params);
348 Ok(Box::new(BufferedIndexScanSource::node_text(
349 upstream, ctx, op,
350 )))
351 }
352
353 PhysicalOp::NodeByPointScan(op @ NodeByPointScanExec { input, .. }) => {
354 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
355 let ctx = StreamCtx::new(storage, params);
356 Ok(Box::new(BufferedIndexScanSource::node_point(
357 upstream, ctx, op,
358 )))
359 }
360
361 PhysicalOp::RelByPropertyRangeScan(op @ RelByPropertyRangeScanExec { input, .. }) => {
362 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
363 let ctx = StreamCtx::new(storage, params);
364 Ok(Box::new(BufferedIndexScanSource::rel_range(
365 upstream, ctx, op,
366 )))
367 }
368
369 PhysicalOp::RelByTextScan(op @ RelByTextScanExec { input, .. }) => {
370 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
371 let ctx = StreamCtx::new(storage, params);
372 Ok(Box::new(BufferedIndexScanSource::rel_text(
373 upstream, ctx, op,
374 )))
375 }
376
377 PhysicalOp::RelByPointScan(op @ RelByPointScanExec { input, .. }) => {
378 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
379 let ctx = StreamCtx::new(storage, params);
380 Ok(Box::new(BufferedIndexScanSource::rel_point(
381 upstream, ctx, op,
382 )))
383 }
384
385 PhysicalOp::Expand(ExpandExec {
386 input,
387 src,
388 rel,
389 dst,
390 types,
391 direction,
392 rel_properties,
393 range,
394 }) => {
395 let upstream =
396 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
397 let ctx = StreamCtx::new(storage, params);
398 match range.as_ref() {
399 Some(range) => Ok(Box::new(VariableLengthExpandSource::new(
400 upstream, ctx, *src, *rel, *dst, types, *direction, range,
401 ))),
402 None => Ok(Box::new(ExpandSource::new(
403 upstream,
404 ctx,
405 *src,
406 *rel,
407 *dst,
408 types,
409 *direction,
410 rel_properties.as_ref(),
411 ))),
412 }
413 }
414
415 PhysicalOp::Filter(FilterExec { input, predicate }) => {
416 let upstream =
417 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
418 let ctx = StreamCtx::new(storage, params);
419 Ok(Box::new(FilterSource::new(upstream, ctx, predicate)))
420 }
421
422 PhysicalOp::Projection(ProjectionExec {
423 input,
424 distinct,
425 items,
426 include_existing,
427 }) => {
428 let upstream =
429 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
430 let ctx = StreamCtx::new(storage, params);
431 let proj: Box<dyn RowSource + 'a> = Box::new(ProjectionSource::new(
432 upstream,
433 ctx,
434 items,
435 *include_existing,
436 ));
437 if *distinct {
438 Ok(Box::new(DistinctSource::new(proj)))
439 } else {
440 Ok(proj)
441 }
442 }
443
444 PhysicalOp::Unwind(UnwindExec { input, expr, alias }) => {
445 let upstream =
446 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
447 let ctx = StreamCtx::new(storage, params);
448 Ok(Box::new(UnwindSource::new(upstream, ctx, expr, *alias)))
449 }
450
451 PhysicalOp::Limit(LimitExec { input, skip, limit }) => {
452 let upstream =
453 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
454 let ctx = StreamCtx::new(storage, params);
457 let eval_ctx = ctx.eval_ctx();
458 let scratch = Row::new();
459 let skip_n = skip
460 .as_ref()
461 .and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
462 .unwrap_or(0)
463 .max(0) as usize;
464 let limit_n = limit
465 .as_ref()
466 .and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
467 .map(|n| n.max(0) as usize);
468 Ok(Box::new(LimitSource::new(upstream, skip_n, limit_n)))
469 }
470
471 PhysicalOp::Sort(SortExec {
472 input,
473 items,
474 top_k,
475 }) => {
476 let upstream =
477 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
478 let ctx = StreamCtx::new(storage, params);
479 Ok(Box::new(SortSource::new_with_top_k(
480 upstream, ctx, items, *top_k,
481 )))
482 }
483
484 PhysicalOp::HashAggregation(
485 agg @ HashAggregationExec {
486 input,
487 group_by,
488 aggregates,
489 },
490 ) => {
491 if seed.is_none() {
495 if let Some(rows) =
496 crate::executor::count_all_scan_aggregation_rows(storage, plan, agg)
497 {
498 return Ok(Box::new(BufferedRowSource::new(rows)));
499 }
500 }
501 let upstream =
502 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
503 let ctx = StreamCtx::new(storage, params);
504 Ok(Box::new(HashAggregationSource::new(
505 upstream, ctx, group_by, aggregates,
506 )))
507 }
508
509 PhysicalOp::OptionalMatch(OptionalMatchExec {
510 input,
511 inner,
512 new_vars,
513 }) => {
514 let upstream =
515 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
516 let ctx = StreamCtx::new(storage, params);
517 Ok(Box::new(OptionalMatchSource::new(
518 upstream, ctx, plan, *inner, new_vars,
519 )))
520 }
521
522 PhysicalOp::CallSubquery(CallSubqueryExec {
523 input,
524 inner,
525 new_vars,
526 }) => {
527 let upstream = build_streaming_dispatch(plan, *input, storage, params.clone(), seed)?;
528 Ok(Box::new(CallSubquerySource::new(
529 upstream, plan, *inner, storage, params, new_vars,
530 )))
531 }
532
533 PhysicalOp::PathBuild(PathBuildExec {
534 input,
535 output,
536 node_vars,
537 rel_vars,
538 shortest_path_all,
539 }) => {
540 let upstream =
541 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
542 let ctx = StreamCtx::new(storage, params);
543 Ok(Box::new(PathBuildSource::new(
544 upstream,
545 ctx,
546 *output,
547 node_vars,
548 rel_vars,
549 *shortest_path_all,
550 )))
551 }
552
553 _ => Err(ExecutorError::RuntimeError(format!(
556 "non-streaming op reached streaming branch: {op:?}"
557 ))),
558 }
559}
560
561fn open_input<'a, S: GraphStorage + 'a>(
565 plan: &'a PhysicalPlan,
566 input: Option<PhysicalNodeId>,
567 storage: &'a S,
568 params: Arc<BTreeMap<String, LoraValue>>,
569 seed: Option<Row>,
570) -> ExecResult<Box<dyn RowSource + 'a>> {
571 match input {
572 Some(input) => build_streaming_dispatch(plan, input, storage, params, seed),
573 None => match seed {
574 Some(seed_row) => Ok(Box::new(BufferedRowSource::new(vec![seed_row]))),
575 None => Ok(Box::new(ArgumentSource::new())),
576 },
577 }
578}
579
580fn build_buffered_subtree<'a, S: GraphStorage + 'a>(
587 plan: &'a PhysicalPlan,
588 node_id: PhysicalNodeId,
589 storage: &'a S,
590 params: &Arc<BTreeMap<String, LoraValue>>,
591) -> ExecResult<Box<dyn RowSource + 'a>> {
592 let executor = Executor::with_deadline(
596 ExecutionContext {
597 storage,
598 params: (**params).clone(),
599 },
600 crate::cancel::active_deadline(),
601 );
602 let rows = executor.execute_subtree(plan, node_id)?;
603 Ok(Box::new(BufferedRowSource::new(rows)))
604}
605
606pub struct PullExecutor<'a, S: GraphStorage> {
612 storage: &'a S,
613 params: BTreeMap<String, LoraValue>,
614}
615
616impl<'a, S: GraphStorage> PullExecutor<'a, S> {
617 pub fn new(storage: &'a S, params: BTreeMap<String, LoraValue>) -> Self {
618 Self { storage, params }
619 }
620
621 pub fn open_compiled(self, compiled: &'a CompiledQuery) -> ExecResult<Box<dyn RowSource + 'a>>
630 where
631 S: 'a,
632 {
633 clear_eval_error();
634 compiled_to_streaming(compiled, self.storage, self.params)
635 }
636}
637
638pub fn collect_compiled<'a, S: GraphStorage + 'a>(
641 storage: &'a S,
642 params: BTreeMap<String, LoraValue>,
643 compiled: &'a CompiledQuery,
644) -> ExecResult<Vec<Row>> {
645 let mut cursor = PullExecutor::new(storage, params).open_compiled(compiled)?;
646 drain(cursor.as_mut())
647}
648
649pub fn collect_compiled_with_deadline<'a, S: GraphStorage + 'a>(
656 storage: &'a S,
657 params: BTreeMap<String, LoraValue>,
658 compiled: &'a CompiledQuery,
659 deadline: Option<web_time::Instant>,
660) -> ExecResult<Vec<Row>> {
661 let _deadline_scope = crate::cancel::DeadlineScope::enter(deadline);
662 if let Some(deadline) = deadline {
663 if crate::cancel::deadline_reached(deadline) {
664 return Err(ExecutorError::QueryTimeout);
665 }
666 }
667 let mut cursor = PullExecutor::new(storage, params).open_compiled(compiled)?;
668 drain(cursor.as_mut())
669}