1use std::collections::BTreeMap;
9use std::sync::Arc;
10
11use lora_compiler::physical::{
12 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::expand::{ExpandSource, VariableLengthExpandSource};
29use super::filter::FilterSource;
30use super::optional::OptionalMatchSource;
31use super::path::PathBuildSource;
32use super::projection::{DistinctSource, ProjectionSource, UnwindSource};
33use super::scan::{
34 BufferedIndexScanSource, NodeByLabelScanSource, NodeByPropertyScanSource, NodeScanSource,
35};
36use super::sort::{LimitSource, SortSource};
37use super::union::UnionSource;
38use super::{drain, ArgumentSource, BufferedRowSource, HydratingSource, RowSource, StreamCtx};
39
40pub(super) fn compiled_to_streaming<'a, S: GraphStorage + 'a>(
56 compiled: &'a CompiledQuery,
57 storage: &'a S,
58 params: BTreeMap<String, LoraValue>,
59) -> ExecResult<Box<dyn RowSource + 'a>> {
60 let params = Arc::new(params);
61
62 if compiled.unions.is_empty() {
63 let plan = &compiled.physical;
64 let inner = build_streaming(plan, plan.root, storage, params)?;
65 if !plan_may_need_hydration(plan) {
66 return Ok(inner);
67 }
68 return Ok(Box::new(HydratingSource::new(inner, storage)));
69 }
70
71 let mut branches: Vec<Box<dyn RowSource + 'a>> = Vec::with_capacity(compiled.unions.len() + 1);
72
73 let head_inner = build_streaming(
74 &compiled.physical,
75 compiled.physical.root,
76 storage,
77 params.clone(),
78 )?;
79 if plan_may_need_hydration(&compiled.physical) {
80 branches.push(Box::new(HydratingSource::new(head_inner, storage)));
81 } else {
82 branches.push(head_inner);
83 }
84
85 let mut needs_dedup = false;
86 for branch in &compiled.unions {
87 let inner = build_streaming(
88 &branch.physical,
89 branch.physical.root,
90 storage,
91 params.clone(),
92 )?;
93 if plan_may_need_hydration(&branch.physical) {
94 branches.push(Box::new(HydratingSource::new(inner, storage)));
95 } else {
96 branches.push(inner);
97 }
98 if !branch.all {
99 needs_dedup = true;
100 }
101 }
102
103 Ok(Box::new(UnionSource::new(branches, needs_dedup)))
104}
105
106pub(super) fn is_streaming_op(op: &PhysicalOp) -> bool {
114 match op {
115 PhysicalOp::Argument(_)
116 | PhysicalOp::NodeScan(_)
117 | PhysicalOp::NodeByLabelScan(_)
118 | PhysicalOp::NodeByPropertyScan(_)
119 | PhysicalOp::NodeByPropertyRangeScan(_)
120 | PhysicalOp::NodeByTextScan(_)
121 | PhysicalOp::NodeByPointScan(_)
122 | PhysicalOp::RelByPropertyRangeScan(_)
123 | PhysicalOp::RelByTextScan(_)
124 | PhysicalOp::RelByPointScan(_)
125 | PhysicalOp::Filter(_)
126 | PhysicalOp::Unwind(_)
127 | PhysicalOp::Limit(_)
128 | PhysicalOp::Sort(_)
135 | PhysicalOp::HashAggregation(_)
136 | PhysicalOp::OptionalMatch(_)
137 | PhysicalOp::PathBuild(_)
138 | PhysicalOp::Projection(_) => true,
142 PhysicalOp::Expand(_) => true,
145 _ => false,
146 }
147}
148
149pub(crate) fn write_op_input(
154 plan: &PhysicalPlan,
155 node_id: PhysicalNodeId,
156) -> Option<PhysicalNodeId> {
157 match &plan.nodes[node_id] {
158 PhysicalOp::Create(o) => Some(o.input),
159 PhysicalOp::Set(o) => Some(o.input),
160 PhysicalOp::Delete(o) => Some(o.input),
161 PhysicalOp::Remove(o) => Some(o.input),
162 PhysicalOp::Merge(o) => Some(o.input),
163 _ => None,
164 }
165}
166
167pub(crate) fn subtree_is_fully_streaming(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
174 let op = &plan.nodes[node_id];
175 if !is_streaming_op(op) {
176 return false;
177 }
178 let child = match op {
179 PhysicalOp::Argument(_) => return true,
180 PhysicalOp::NodeScan(o) => o.input,
181 PhysicalOp::NodeByLabelScan(o) => o.input,
182 PhysicalOp::NodeByPropertyScan(o) => o.input,
183 PhysicalOp::NodeByPropertyRangeScan(o) => o.input,
184 PhysicalOp::NodeByTextScan(o) => o.input,
185 PhysicalOp::NodeByPointScan(o) => o.input,
186 PhysicalOp::RelByPropertyRangeScan(o) => o.input,
187 PhysicalOp::RelByTextScan(o) => o.input,
188 PhysicalOp::RelByPointScan(o) => o.input,
189 PhysicalOp::Filter(o) => Some(o.input),
190 PhysicalOp::Unwind(o) => Some(o.input),
191 PhysicalOp::Limit(o) => Some(o.input),
192 PhysicalOp::Expand(o) => Some(o.input),
193 PhysicalOp::Projection(o) => Some(o.input),
194 PhysicalOp::Sort(o) => Some(o.input),
195 PhysicalOp::HashAggregation(o) => Some(o.input),
196 PhysicalOp::OptionalMatch(o) => Some(o.input),
197 PhysicalOp::PathBuild(o) => Some(o.input),
198 _ => return false,
200 };
201 match child {
202 None => true,
203 Some(c) => subtree_is_fully_streaming(plan, c),
204 }
205}
206
207pub(crate) fn build_streaming<'a, S: GraphStorage + 'a>(
208 plan: &'a PhysicalPlan,
209 node_id: PhysicalNodeId,
210 storage: &'a S,
211 params: Arc<BTreeMap<String, LoraValue>>,
212) -> ExecResult<Box<dyn RowSource + 'a>> {
213 build_streaming_inner(plan, node_id, storage, params).map(|src| wrap_metered(node_id, src))
214}
215
216fn build_streaming_inner<'a, S: GraphStorage + 'a>(
217 plan: &'a PhysicalPlan,
218 node_id: PhysicalNodeId,
219 storage: &'a S,
220 params: Arc<BTreeMap<String, LoraValue>>,
221) -> ExecResult<Box<dyn RowSource + 'a>> {
222 let op = &plan.nodes[node_id];
223
224 if !is_streaming_op(op) {
225 return build_buffered_subtree(plan, node_id, storage, ¶ms);
226 }
227
228 match op {
229 PhysicalOp::Argument(_) => Ok(Box::new(ArgumentSource::new())),
230
231 PhysicalOp::NodeScan(NodeScanExec { input, var }) => {
232 let upstream = open_input(plan, *input, storage, params.clone())?;
233 Ok(Box::new(NodeScanSource::new(upstream, storage, *var)))
234 }
235
236 PhysicalOp::NodeByLabelScan(NodeByLabelScanExec { input, var, labels }) => {
237 let upstream = open_input(plan, *input, storage, params.clone())?;
238 Ok(Box::new(NodeByLabelScanSource::new(
239 upstream, storage, *var, labels,
240 )))
241 }
242
243 PhysicalOp::NodeByPropertyScan(NodeByPropertyScanExec {
244 input,
245 var,
246 labels,
247 key,
248 value,
249 }) => {
250 let upstream = open_input(plan, *input, storage, params.clone())?;
251 let ctx = StreamCtx::new(storage, params);
252 Ok(Box::new(NodeByPropertyScanSource::new(
253 upstream, ctx, *var, labels, key, value,
254 )))
255 }
256
257 PhysicalOp::NodeByPropertyRangeScan(op @ NodeByPropertyRangeScanExec { input, .. }) => {
258 let upstream = open_input(plan, *input, storage, params.clone())?;
259 let ctx = StreamCtx::new(storage, params);
260 Ok(Box::new(BufferedIndexScanSource::node_range(
261 upstream, ctx, op,
262 )))
263 }
264
265 PhysicalOp::NodeByTextScan(op @ NodeByTextScanExec { input, .. }) => {
266 let upstream = open_input(plan, *input, storage, params.clone())?;
267 let ctx = StreamCtx::new(storage, params);
268 Ok(Box::new(BufferedIndexScanSource::node_text(
269 upstream, ctx, op,
270 )))
271 }
272
273 PhysicalOp::NodeByPointScan(op @ NodeByPointScanExec { input, .. }) => {
274 let upstream = open_input(plan, *input, storage, params.clone())?;
275 let ctx = StreamCtx::new(storage, params);
276 Ok(Box::new(BufferedIndexScanSource::node_point(
277 upstream, ctx, op,
278 )))
279 }
280
281 PhysicalOp::RelByPropertyRangeScan(op @ RelByPropertyRangeScanExec { input, .. }) => {
282 let upstream = open_input(plan, *input, storage, params.clone())?;
283 let ctx = StreamCtx::new(storage, params);
284 Ok(Box::new(BufferedIndexScanSource::rel_range(
285 upstream, ctx, op,
286 )))
287 }
288
289 PhysicalOp::RelByTextScan(op @ RelByTextScanExec { input, .. }) => {
290 let upstream = open_input(plan, *input, storage, params.clone())?;
291 let ctx = StreamCtx::new(storage, params);
292 Ok(Box::new(BufferedIndexScanSource::rel_text(
293 upstream, ctx, op,
294 )))
295 }
296
297 PhysicalOp::RelByPointScan(op @ RelByPointScanExec { input, .. }) => {
298 let upstream = open_input(plan, *input, storage, params.clone())?;
299 let ctx = StreamCtx::new(storage, params);
300 Ok(Box::new(BufferedIndexScanSource::rel_point(
301 upstream, ctx, op,
302 )))
303 }
304
305 PhysicalOp::Expand(ExpandExec {
306 input,
307 src,
308 rel,
309 dst,
310 types,
311 direction,
312 rel_properties,
313 range,
314 }) => {
315 let upstream = build_streaming(plan, *input, storage, params.clone())?;
316 let ctx = StreamCtx::new(storage, params);
317 match range.as_ref() {
318 Some(range) => Ok(Box::new(VariableLengthExpandSource::new(
319 upstream, ctx, *src, *rel, *dst, types, *direction, range,
320 ))),
321 None => Ok(Box::new(ExpandSource::new(
322 upstream,
323 ctx,
324 *src,
325 *rel,
326 *dst,
327 types,
328 *direction,
329 rel_properties.as_ref(),
330 ))),
331 }
332 }
333
334 PhysicalOp::Filter(FilterExec { input, predicate }) => {
335 let upstream = build_streaming(plan, *input, storage, params.clone())?;
336 let ctx = StreamCtx::new(storage, params);
337 Ok(Box::new(FilterSource::new(upstream, ctx, predicate)))
338 }
339
340 PhysicalOp::Projection(ProjectionExec {
341 input,
342 distinct,
343 items,
344 include_existing,
345 }) => {
346 let upstream = build_streaming(plan, *input, storage, params.clone())?;
347 let ctx = StreamCtx::new(storage, params);
348 let proj: Box<dyn RowSource + 'a> = Box::new(ProjectionSource::new(
349 upstream,
350 ctx,
351 items,
352 *include_existing,
353 ));
354 if *distinct {
355 Ok(Box::new(DistinctSource::new(proj)))
356 } else {
357 Ok(proj)
358 }
359 }
360
361 PhysicalOp::Unwind(UnwindExec { input, expr, alias }) => {
362 let upstream = build_streaming(plan, *input, storage, params.clone())?;
363 let ctx = StreamCtx::new(storage, params);
364 Ok(Box::new(UnwindSource::new(upstream, ctx, expr, *alias)))
365 }
366
367 PhysicalOp::Limit(LimitExec { input, skip, limit }) => {
368 let upstream = build_streaming(plan, *input, storage, params.clone())?;
369 let ctx = StreamCtx::new(storage, params);
372 let eval_ctx = ctx.eval_ctx();
373 let scratch = Row::new();
374 let skip_n = skip
375 .as_ref()
376 .and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
377 .unwrap_or(0)
378 .max(0) as usize;
379 let limit_n = limit
380 .as_ref()
381 .and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
382 .map(|n| n.max(0) as usize);
383 Ok(Box::new(LimitSource::new(upstream, skip_n, limit_n)))
384 }
385
386 PhysicalOp::Sort(SortExec {
387 input,
388 items,
389 top_k,
390 }) => {
391 let upstream = build_streaming(plan, *input, storage, params.clone())?;
392 let ctx = StreamCtx::new(storage, params);
393 Ok(Box::new(SortSource::new_with_top_k(
394 upstream, ctx, items, *top_k,
395 )))
396 }
397
398 PhysicalOp::HashAggregation(HashAggregationExec {
399 input,
400 group_by,
401 aggregates,
402 }) => {
403 let upstream = build_streaming(plan, *input, storage, params.clone())?;
404 let ctx = StreamCtx::new(storage, params);
405 Ok(Box::new(HashAggregationSource::new(
406 upstream, ctx, group_by, aggregates,
407 )))
408 }
409
410 PhysicalOp::OptionalMatch(OptionalMatchExec {
411 input,
412 inner,
413 new_vars,
414 }) => {
415 let upstream = build_streaming(plan, *input, storage, params.clone())?;
416 let ctx = StreamCtx::new(storage, params);
417 Ok(Box::new(OptionalMatchSource::new(
418 upstream, ctx, plan, *inner, new_vars,
419 )))
420 }
421
422 PhysicalOp::PathBuild(PathBuildExec {
423 input,
424 output,
425 node_vars,
426 rel_vars,
427 shortest_path_all,
428 }) => {
429 let upstream = build_streaming(plan, *input, storage, params.clone())?;
430 let ctx = StreamCtx::new(storage, params);
431 Ok(Box::new(PathBuildSource::new(
432 upstream,
433 ctx,
434 *output,
435 node_vars,
436 rel_vars,
437 *shortest_path_all,
438 )))
439 }
440
441 _ => Err(ExecutorError::RuntimeError(format!(
444 "non-streaming op reached streaming branch: {op:?}"
445 ))),
446 }
447}
448
449fn open_input<'a, S: GraphStorage + 'a>(
453 plan: &'a PhysicalPlan,
454 input: Option<PhysicalNodeId>,
455 storage: &'a S,
456 params: Arc<BTreeMap<String, LoraValue>>,
457) -> ExecResult<Box<dyn RowSource + 'a>> {
458 match input {
459 Some(input) => build_streaming(plan, input, storage, params),
460 None => Ok(Box::new(ArgumentSource::new())),
461 }
462}
463
464fn build_buffered_subtree<'a, S: GraphStorage + 'a>(
471 plan: &'a PhysicalPlan,
472 node_id: PhysicalNodeId,
473 storage: &'a S,
474 params: &Arc<BTreeMap<String, LoraValue>>,
475) -> ExecResult<Box<dyn RowSource + 'a>> {
476 let executor = Executor::new(ExecutionContext {
480 storage,
481 params: (**params).clone(),
482 });
483 let rows = executor.execute_subtree(plan, node_id)?;
484 Ok(Box::new(BufferedRowSource::new(rows)))
485}
486
487pub struct PullExecutor<'a, S: GraphStorage> {
493 storage: &'a S,
494 params: BTreeMap<String, LoraValue>,
495}
496
497impl<'a, S: GraphStorage> PullExecutor<'a, S> {
498 pub fn new(storage: &'a S, params: BTreeMap<String, LoraValue>) -> Self {
499 Self { storage, params }
500 }
501
502 pub fn open_compiled(self, compiled: &'a CompiledQuery) -> ExecResult<Box<dyn RowSource + 'a>>
511 where
512 S: 'a,
513 {
514 clear_eval_error();
515 compiled_to_streaming(compiled, self.storage, self.params)
516 }
517}
518
519pub fn collect_compiled<'a, S: GraphStorage + 'a>(
522 storage: &'a S,
523 params: BTreeMap<String, LoraValue>,
524 compiled: &'a CompiledQuery,
525) -> ExecResult<Vec<Row>> {
526 let mut cursor = PullExecutor::new(storage, params).open_compiled(compiled)?;
527 drain(cursor.as_mut())
528}