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) => Some(o.input),
200 PhysicalOp::PathBuild(o) => Some(o.input),
201 _ => return false,
203 };
204 match child {
205 None => true,
206 Some(c) => subtree_is_fully_streaming(plan, c),
207 }
208}
209
210pub(crate) fn build_streaming<'a, S: GraphStorage + 'a>(
211 plan: &'a PhysicalPlan,
212 node_id: PhysicalNodeId,
213 storage: &'a S,
214 params: Arc<BTreeMap<String, LoraValue>>,
215) -> ExecResult<Box<dyn RowSource + 'a>> {
216 build_streaming_dispatch(plan, node_id, storage, params, None)
217}
218
219pub(crate) fn build_streaming_seeded<'a, S: GraphStorage + 'a>(
223 plan: &'a PhysicalPlan,
224 node_id: PhysicalNodeId,
225 storage: &'a S,
226 params: Arc<BTreeMap<String, LoraValue>>,
227 seed: Row,
228) -> ExecResult<Box<dyn RowSource + 'a>> {
229 build_streaming_dispatch(plan, node_id, storage, params, Some(seed))
230}
231
232fn build_streaming_dispatch<'a, S: GraphStorage + 'a>(
233 plan: &'a PhysicalPlan,
234 node_id: PhysicalNodeId,
235 storage: &'a S,
236 params: Arc<BTreeMap<String, LoraValue>>,
237 seed: Option<Row>,
238) -> ExecResult<Box<dyn RowSource + 'a>> {
239 build_streaming_inner(plan, node_id, storage, params, seed)
240 .map(|src| wrap_metered(node_id, src))
241}
242
243fn build_streaming_inner<'a, S: GraphStorage + 'a>(
244 plan: &'a PhysicalPlan,
245 node_id: PhysicalNodeId,
246 storage: &'a S,
247 params: Arc<BTreeMap<String, LoraValue>>,
248 seed: Option<Row>,
249) -> ExecResult<Box<dyn RowSource + 'a>> {
250 let op = &plan.nodes[node_id];
251
252 if !is_streaming_op(op) {
253 return build_buffered_subtree(plan, node_id, storage, ¶ms);
254 }
255
256 match op {
257 PhysicalOp::Argument(_) => match seed {
258 Some(seed_row) => Ok(Box::new(BufferedRowSource::new(vec![seed_row]))),
259 None => Ok(Box::new(ArgumentSource::new())),
260 },
261
262 PhysicalOp::NodeScan(NodeScanExec { input, var }) => {
263 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
264 Ok(Box::new(NodeScanSource::new(upstream, storage, *var)))
265 }
266
267 PhysicalOp::NodeByLabelScan(NodeByLabelScanExec { input, var, labels }) => {
268 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
269 Ok(Box::new(NodeByLabelScanSource::new(
270 upstream, storage, *var, labels,
271 )))
272 }
273
274 PhysicalOp::NodeByPropertyScan(NodeByPropertyScanExec {
275 input,
276 var,
277 labels,
278 key,
279 value,
280 }) => {
281 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
282 let ctx = StreamCtx::new(storage, params);
283 Ok(Box::new(NodeByPropertyScanSource::new(
284 upstream, ctx, *var, labels, key, value,
285 )))
286 }
287
288 PhysicalOp::NodeByPropertyRangeScan(op @ NodeByPropertyRangeScanExec { input, .. }) => {
289 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
290 let ctx = StreamCtx::new(storage, params);
291 Ok(Box::new(BufferedIndexScanSource::node_range(
292 upstream, ctx, op,
293 )))
294 }
295
296 PhysicalOp::NodeByTextScan(op @ NodeByTextScanExec { input, .. }) => {
297 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
298 let ctx = StreamCtx::new(storage, params);
299 Ok(Box::new(BufferedIndexScanSource::node_text(
300 upstream, ctx, op,
301 )))
302 }
303
304 PhysicalOp::NodeByPointScan(op @ NodeByPointScanExec { input, .. }) => {
305 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
306 let ctx = StreamCtx::new(storage, params);
307 Ok(Box::new(BufferedIndexScanSource::node_point(
308 upstream, ctx, op,
309 )))
310 }
311
312 PhysicalOp::RelByPropertyRangeScan(op @ RelByPropertyRangeScanExec { input, .. }) => {
313 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
314 let ctx = StreamCtx::new(storage, params);
315 Ok(Box::new(BufferedIndexScanSource::rel_range(
316 upstream, ctx, op,
317 )))
318 }
319
320 PhysicalOp::RelByTextScan(op @ RelByTextScanExec { input, .. }) => {
321 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
322 let ctx = StreamCtx::new(storage, params);
323 Ok(Box::new(BufferedIndexScanSource::rel_text(
324 upstream, ctx, op,
325 )))
326 }
327
328 PhysicalOp::RelByPointScan(op @ RelByPointScanExec { input, .. }) => {
329 let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
330 let ctx = StreamCtx::new(storage, params);
331 Ok(Box::new(BufferedIndexScanSource::rel_point(
332 upstream, ctx, op,
333 )))
334 }
335
336 PhysicalOp::Expand(ExpandExec {
337 input,
338 src,
339 rel,
340 dst,
341 types,
342 direction,
343 rel_properties,
344 range,
345 }) => {
346 let upstream =
347 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
348 let ctx = StreamCtx::new(storage, params);
349 match range.as_ref() {
350 Some(range) => Ok(Box::new(VariableLengthExpandSource::new(
351 upstream, ctx, *src, *rel, *dst, types, *direction, range,
352 ))),
353 None => Ok(Box::new(ExpandSource::new(
354 upstream,
355 ctx,
356 *src,
357 *rel,
358 *dst,
359 types,
360 *direction,
361 rel_properties.as_ref(),
362 ))),
363 }
364 }
365
366 PhysicalOp::Filter(FilterExec { input, predicate }) => {
367 let upstream =
368 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
369 let ctx = StreamCtx::new(storage, params);
370 Ok(Box::new(FilterSource::new(upstream, ctx, predicate)))
371 }
372
373 PhysicalOp::Projection(ProjectionExec {
374 input,
375 distinct,
376 items,
377 include_existing,
378 }) => {
379 let upstream =
380 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
381 let ctx = StreamCtx::new(storage, params);
382 let proj: Box<dyn RowSource + 'a> = Box::new(ProjectionSource::new(
383 upstream,
384 ctx,
385 items,
386 *include_existing,
387 ));
388 if *distinct {
389 Ok(Box::new(DistinctSource::new(proj)))
390 } else {
391 Ok(proj)
392 }
393 }
394
395 PhysicalOp::Unwind(UnwindExec { input, expr, alias }) => {
396 let upstream =
397 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
398 let ctx = StreamCtx::new(storage, params);
399 Ok(Box::new(UnwindSource::new(upstream, ctx, expr, *alias)))
400 }
401
402 PhysicalOp::Limit(LimitExec { input, skip, limit }) => {
403 let upstream =
404 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
405 let ctx = StreamCtx::new(storage, params);
408 let eval_ctx = ctx.eval_ctx();
409 let scratch = Row::new();
410 let skip_n = skip
411 .as_ref()
412 .and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
413 .unwrap_or(0)
414 .max(0) as usize;
415 let limit_n = limit
416 .as_ref()
417 .and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
418 .map(|n| n.max(0) as usize);
419 Ok(Box::new(LimitSource::new(upstream, skip_n, limit_n)))
420 }
421
422 PhysicalOp::Sort(SortExec {
423 input,
424 items,
425 top_k,
426 }) => {
427 let upstream =
428 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
429 let ctx = StreamCtx::new(storage, params);
430 Ok(Box::new(SortSource::new_with_top_k(
431 upstream, ctx, items, *top_k,
432 )))
433 }
434
435 PhysicalOp::HashAggregation(HashAggregationExec {
436 input,
437 group_by,
438 aggregates,
439 }) => {
440 let upstream =
441 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
442 let ctx = StreamCtx::new(storage, params);
443 Ok(Box::new(HashAggregationSource::new(
444 upstream, ctx, group_by, aggregates,
445 )))
446 }
447
448 PhysicalOp::OptionalMatch(OptionalMatchExec {
449 input,
450 inner,
451 new_vars,
452 }) => {
453 let upstream =
454 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
455 let ctx = StreamCtx::new(storage, params);
456 Ok(Box::new(OptionalMatchSource::new(
457 upstream, ctx, plan, *inner, new_vars,
458 )))
459 }
460
461 PhysicalOp::CallSubquery(CallSubqueryExec {
462 input,
463 inner,
464 new_vars,
465 }) => {
466 let upstream = build_streaming_dispatch(plan, *input, storage, params.clone(), seed)?;
467 Ok(Box::new(CallSubquerySource::new(
468 upstream, plan, *inner, storage, params, new_vars,
469 )))
470 }
471
472 PhysicalOp::PathBuild(PathBuildExec {
473 input,
474 output,
475 node_vars,
476 rel_vars,
477 shortest_path_all,
478 }) => {
479 let upstream =
480 build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
481 let ctx = StreamCtx::new(storage, params);
482 Ok(Box::new(PathBuildSource::new(
483 upstream,
484 ctx,
485 *output,
486 node_vars,
487 rel_vars,
488 *shortest_path_all,
489 )))
490 }
491
492 _ => Err(ExecutorError::RuntimeError(format!(
495 "non-streaming op reached streaming branch: {op:?}"
496 ))),
497 }
498}
499
500fn open_input<'a, S: GraphStorage + 'a>(
504 plan: &'a PhysicalPlan,
505 input: Option<PhysicalNodeId>,
506 storage: &'a S,
507 params: Arc<BTreeMap<String, LoraValue>>,
508 seed: Option<Row>,
509) -> ExecResult<Box<dyn RowSource + 'a>> {
510 match input {
511 Some(input) => build_streaming_dispatch(plan, input, storage, params, seed),
512 None => match seed {
513 Some(seed_row) => Ok(Box::new(BufferedRowSource::new(vec![seed_row]))),
514 None => Ok(Box::new(ArgumentSource::new())),
515 },
516 }
517}
518
519fn build_buffered_subtree<'a, S: GraphStorage + 'a>(
526 plan: &'a PhysicalPlan,
527 node_id: PhysicalNodeId,
528 storage: &'a S,
529 params: &Arc<BTreeMap<String, LoraValue>>,
530) -> ExecResult<Box<dyn RowSource + 'a>> {
531 let executor = Executor::new(ExecutionContext {
535 storage,
536 params: (**params).clone(),
537 });
538 let rows = executor.execute_subtree(plan, node_id)?;
539 Ok(Box::new(BufferedRowSource::new(rows)))
540}
541
542pub struct PullExecutor<'a, S: GraphStorage> {
548 storage: &'a S,
549 params: BTreeMap<String, LoraValue>,
550}
551
552impl<'a, S: GraphStorage> PullExecutor<'a, S> {
553 pub fn new(storage: &'a S, params: BTreeMap<String, LoraValue>) -> Self {
554 Self { storage, params }
555 }
556
557 pub fn open_compiled(self, compiled: &'a CompiledQuery) -> ExecResult<Box<dyn RowSource + 'a>>
566 where
567 S: 'a,
568 {
569 clear_eval_error();
570 compiled_to_streaming(compiled, self.storage, self.params)
571 }
572}
573
574pub fn collect_compiled<'a, S: GraphStorage + 'a>(
577 storage: &'a S,
578 params: BTreeMap<String, LoraValue>,
579 compiled: &'a CompiledQuery,
580) -> ExecResult<Vec<Row>> {
581 let mut cursor = PullExecutor::new(storage, params).open_compiled(compiled)?;
582 drain(cursor.as_mut())
583}