1use std::collections::{BTreeMap, BTreeSet};
21use std::mem::ManuallyDrop;
22use std::sync::Arc;
23
24use lora_compiler::physical::{
25 ExpandExec, FilterExec, HashAggregationExec, LimitExec, NodeByLabelScanExec,
26 NodeByPropertyScanExec, NodeScanExec, OptionalMatchExec, PathBuildExec, PhysicalNodeId,
27 PhysicalOp, PhysicalPlan, ProjectionExec, SortExec, UnwindExec,
28};
29use lora_compiler::CompiledQuery;
30use lora_store::{GraphStorage, GraphStorageMut};
31
32use crate::errors::{ExecResult, ExecutorError};
33use crate::eval::{clear_eval_error, eval_expr, EvalContext};
34use crate::executor::{
35 hydrate_node_record, hydrate_relationship_record, ExecutionContext, Executor, GroupValueKey,
36 MutableExecutionContext, MutableExecutor,
37};
38use crate::value::{LoraValue, Row};
39
40use super::aggregate::HashAggregationSource;
41use super::expand::{ExpandSource, VariableLengthExpandSource};
42use super::filter::FilterSource;
43use super::optional::OptionalMatchSource;
44use super::path::PathBuildSource;
45use super::projection::{DistinctSource, ProjectionSource, UnwindSource};
46use super::scan::{NodeByLabelScanSource, NodeByPropertyScanSource, NodeScanSource};
47use super::sort::{LimitSource, SortSource};
48use super::union::UnionSource;
49
50pub trait RowSource {
57 fn next_row(&mut self) -> ExecResult<Option<Row>>;
59}
60
61pub fn drain<S: RowSource + ?Sized>(source: &mut S) -> ExecResult<Vec<Row>> {
63 let mut out = Vec::new();
64 while let Some(row) = source.next_row()? {
65 out.push(row);
66 }
67 Ok(out)
68}
69
70#[derive(Clone)]
80pub(super) struct StreamCtx<'a, S: GraphStorage> {
81 pub storage: &'a S,
82 pub params: Arc<BTreeMap<String, LoraValue>>,
83}
84
85impl<'a, S: GraphStorage> StreamCtx<'a, S> {
86 pub(super) fn new(storage: &'a S, params: Arc<BTreeMap<String, LoraValue>>) -> Self {
87 Self { storage, params }
88 }
89
90 pub(super) fn eval_ctx<'b>(&'b self) -> EvalContext<'b, S> {
93 EvalContext {
94 storage: self.storage,
95 params: &self.params,
96 }
97 }
98}
99
100pub struct BufferedRowSource {
108 iter: std::vec::IntoIter<Row>,
109}
110
111impl BufferedRowSource {
112 pub fn new(rows: Vec<Row>) -> Self {
113 Self {
114 iter: rows.into_iter(),
115 }
116 }
117}
118
119impl RowSource for BufferedRowSource {
120 fn next_row(&mut self) -> ExecResult<Option<Row>> {
121 Ok(self.iter.next())
122 }
123}
124
125pub struct ArgumentSource {
132 yielded: bool,
133}
134
135impl ArgumentSource {
136 pub fn new() -> Self {
137 Self { yielded: false }
138 }
139}
140
141impl Default for ArgumentSource {
142 fn default() -> Self {
143 Self::new()
144 }
145}
146
147impl RowSource for ArgumentSource {
148 fn next_row(&mut self) -> ExecResult<Option<Row>> {
149 if self.yielded {
150 Ok(None)
151 } else {
152 self.yielded = true;
153 Ok(Some(Row::new()))
154 }
155 }
156}
157
158pub struct HydratingSource<'a, S: GraphStorage> {
166 upstream: Box<dyn RowSource + 'a>,
167 storage: &'a S,
168}
169
170impl<'a, S: GraphStorage> HydratingSource<'a, S> {
171 pub(super) fn new(upstream: Box<dyn RowSource + 'a>, storage: &'a S) -> Self {
172 Self { upstream, storage }
173 }
174}
175
176impl<'a, S: GraphStorage> RowSource for HydratingSource<'a, S> {
177 fn next_row(&mut self) -> ExecResult<Option<Row>> {
178 match self.upstream.next_row()? {
179 None => Ok(None),
180 Some(row) => {
181 let mut out = Row::new();
182 for (var, name, value) in row.into_iter_named() {
183 out.insert_named(var, name, hydrate_value(value, self.storage));
184 }
185 Ok(Some(out))
186 }
187 }
188 }
189}
190
191pub(super) fn hydrate_value<S: GraphStorage>(value: LoraValue, storage: &S) -> LoraValue {
192 match value {
193 LoraValue::Node(id) => storage
194 .with_node(id, hydrate_node_record)
195 .unwrap_or(LoraValue::Null),
196 LoraValue::Relationship(id) => storage
197 .with_relationship(id, hydrate_relationship_record)
198 .unwrap_or(LoraValue::Null),
199 LoraValue::List(values) => LoraValue::List(
200 values
201 .into_iter()
202 .map(|v| hydrate_value(v, storage))
203 .collect(),
204 ),
205 LoraValue::Map(map) => LoraValue::Map(
206 map.into_iter()
207 .map(|(k, v)| (k, hydrate_value(v, storage)))
208 .collect(),
209 ),
210 other => other,
211 }
212}
213
214pub(super) fn compiled_to_streaming<'a, S: GraphStorage + 'a>(
230 compiled: &'a CompiledQuery,
231 storage: &'a S,
232 params: BTreeMap<String, LoraValue>,
233) -> ExecResult<Box<dyn RowSource + 'a>> {
234 let params = Arc::new(params);
235
236 if compiled.unions.is_empty() {
237 let plan = &compiled.physical;
238 let inner = build_streaming(plan, plan.root, storage, params)?;
239 return Ok(Box::new(HydratingSource::new(inner, storage)));
240 }
241
242 let mut branches: Vec<Box<dyn RowSource + 'a>> = Vec::with_capacity(compiled.unions.len() + 1);
243
244 let head_inner = build_streaming(
245 &compiled.physical,
246 compiled.physical.root,
247 storage,
248 params.clone(),
249 )?;
250 branches.push(Box::new(HydratingSource::new(head_inner, storage)));
251
252 let mut needs_dedup = false;
253 for branch in &compiled.unions {
254 let inner = build_streaming(
255 &branch.physical,
256 branch.physical.root,
257 storage,
258 params.clone(),
259 )?;
260 branches.push(Box::new(HydratingSource::new(inner, storage)));
261 if !branch.all {
262 needs_dedup = true;
263 }
264 }
265
266 Ok(Box::new(UnionSource::new(branches, needs_dedup)))
267}
268
269pub(super) fn is_streaming_op(op: &PhysicalOp) -> bool {
277 match op {
278 PhysicalOp::Argument(_)
279 | PhysicalOp::NodeScan(_)
280 | PhysicalOp::NodeByLabelScan(_)
281 | PhysicalOp::NodeByPropertyScan(_)
282 | PhysicalOp::Filter(_)
283 | PhysicalOp::Unwind(_)
284 | PhysicalOp::Limit(_)
285 | PhysicalOp::Sort(_)
292 | PhysicalOp::HashAggregation(_)
293 | PhysicalOp::OptionalMatch(_)
294 | PhysicalOp::PathBuild(_)
295 | PhysicalOp::Projection(_) => true,
299 PhysicalOp::Expand(_) => true,
302 _ => false,
303 }
304}
305
306pub(super) fn write_op_input(
311 plan: &PhysicalPlan,
312 node_id: PhysicalNodeId,
313) -> Option<PhysicalNodeId> {
314 match &plan.nodes[node_id] {
315 PhysicalOp::Create(o) => Some(o.input),
316 PhysicalOp::Set(o) => Some(o.input),
317 PhysicalOp::Delete(o) => Some(o.input),
318 PhysicalOp::Remove(o) => Some(o.input),
319 PhysicalOp::Merge(o) => Some(o.input),
320 _ => None,
321 }
322}
323
324pub(crate) fn subtree_is_fully_streaming(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
331 let op = &plan.nodes[node_id];
332 if !is_streaming_op(op) {
333 return false;
334 }
335 let child = match op {
336 PhysicalOp::Argument(_) => return true,
337 PhysicalOp::NodeScan(o) => o.input,
338 PhysicalOp::NodeByLabelScan(o) => o.input,
339 PhysicalOp::NodeByPropertyScan(o) => o.input,
340 PhysicalOp::Filter(o) => Some(o.input),
341 PhysicalOp::Unwind(o) => Some(o.input),
342 PhysicalOp::Limit(o) => Some(o.input),
343 PhysicalOp::Expand(o) => Some(o.input),
344 PhysicalOp::Projection(o) => Some(o.input),
345 PhysicalOp::Sort(o) => Some(o.input),
346 PhysicalOp::HashAggregation(o) => Some(o.input),
347 PhysicalOp::OptionalMatch(o) => Some(o.input),
348 PhysicalOp::PathBuild(o) => Some(o.input),
349 _ => return false,
351 };
352 match child {
353 None => true,
354 Some(c) => subtree_is_fully_streaming(plan, c),
355 }
356}
357
358pub(crate) fn build_streaming<'a, S: GraphStorage + 'a>(
359 plan: &'a PhysicalPlan,
360 node_id: PhysicalNodeId,
361 storage: &'a S,
362 params: Arc<BTreeMap<String, LoraValue>>,
363) -> ExecResult<Box<dyn RowSource + 'a>> {
364 let op = &plan.nodes[node_id];
365
366 if !is_streaming_op(op) {
367 return build_buffered_subtree(plan, node_id, storage, ¶ms);
368 }
369
370 match op {
371 PhysicalOp::Argument(_) => Ok(Box::new(ArgumentSource::new())),
372
373 PhysicalOp::NodeScan(NodeScanExec { input, var }) => {
374 let upstream = open_input(plan, *input, storage, params.clone())?;
375 Ok(Box::new(NodeScanSource::new(upstream, storage, *var)))
376 }
377
378 PhysicalOp::NodeByLabelScan(NodeByLabelScanExec { input, var, labels }) => {
379 let upstream = open_input(plan, *input, storage, params.clone())?;
380 Ok(Box::new(NodeByLabelScanSource::new(
381 upstream, storage, *var, labels,
382 )))
383 }
384
385 PhysicalOp::NodeByPropertyScan(NodeByPropertyScanExec {
386 input,
387 var,
388 labels,
389 key,
390 value,
391 }) => {
392 let upstream = open_input(plan, *input, storage, params.clone())?;
393 let ctx = StreamCtx::new(storage, params);
394 Ok(Box::new(NodeByPropertyScanSource::new(
395 upstream, ctx, *var, labels, key, value,
396 )))
397 }
398
399 PhysicalOp::Expand(ExpandExec {
400 input,
401 src,
402 rel,
403 dst,
404 types,
405 direction,
406 rel_properties,
407 range,
408 }) => {
409 let upstream = build_streaming(plan, *input, storage, params.clone())?;
410 let ctx = StreamCtx::new(storage, params);
411 match range.as_ref() {
412 Some(range) => Ok(Box::new(VariableLengthExpandSource::new(
413 upstream, ctx, *src, *rel, *dst, types, *direction, range,
414 ))),
415 None => Ok(Box::new(ExpandSource::new(
416 upstream,
417 ctx,
418 *src,
419 *rel,
420 *dst,
421 types,
422 *direction,
423 rel_properties.as_ref(),
424 ))),
425 }
426 }
427
428 PhysicalOp::Filter(FilterExec { input, predicate }) => {
429 let upstream = build_streaming(plan, *input, storage, params.clone())?;
430 let ctx = StreamCtx::new(storage, params);
431 Ok(Box::new(FilterSource::new(upstream, ctx, predicate)))
432 }
433
434 PhysicalOp::Projection(ProjectionExec {
435 input,
436 distinct,
437 items,
438 include_existing,
439 }) => {
440 let upstream = build_streaming(plan, *input, storage, params.clone())?;
441 let ctx = StreamCtx::new(storage, params);
442 let proj: Box<dyn RowSource + 'a> = Box::new(ProjectionSource::new(
443 upstream,
444 ctx,
445 items,
446 *include_existing,
447 ));
448 if *distinct {
449 Ok(Box::new(DistinctSource::new(proj)))
450 } else {
451 Ok(proj)
452 }
453 }
454
455 PhysicalOp::Unwind(UnwindExec { input, expr, alias }) => {
456 let upstream = build_streaming(plan, *input, storage, params.clone())?;
457 let ctx = StreamCtx::new(storage, params);
458 Ok(Box::new(UnwindSource::new(upstream, ctx, expr, *alias)))
459 }
460
461 PhysicalOp::Limit(LimitExec { input, skip, limit }) => {
462 let upstream = build_streaming(plan, *input, storage, params.clone())?;
463 let ctx = StreamCtx::new(storage, params);
466 let eval_ctx = ctx.eval_ctx();
467 let scratch = Row::new();
468 let skip_n = skip
469 .as_ref()
470 .and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
471 .unwrap_or(0)
472 .max(0) as usize;
473 let limit_n = limit
474 .as_ref()
475 .and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
476 .map(|n| n.max(0) as usize);
477 Ok(Box::new(LimitSource::new(upstream, skip_n, limit_n)))
478 }
479
480 PhysicalOp::Sort(SortExec { input, items }) => {
481 let upstream = build_streaming(plan, *input, storage, params.clone())?;
482 let ctx = StreamCtx::new(storage, params);
483 Ok(Box::new(SortSource::new(upstream, ctx, items)))
484 }
485
486 PhysicalOp::HashAggregation(HashAggregationExec {
487 input,
488 group_by,
489 aggregates,
490 }) => {
491 let upstream = build_streaming(plan, *input, storage, params.clone())?;
492 let ctx = StreamCtx::new(storage, params);
493 Ok(Box::new(HashAggregationSource::new(
494 upstream, ctx, group_by, aggregates,
495 )))
496 }
497
498 PhysicalOp::OptionalMatch(OptionalMatchExec {
499 input,
500 inner,
501 new_vars,
502 }) => {
503 let upstream = build_streaming(plan, *input, storage, params.clone())?;
504 let ctx = StreamCtx::new(storage, params);
505 Ok(Box::new(OptionalMatchSource::new(
506 upstream, ctx, plan, *inner, new_vars,
507 )))
508 }
509
510 PhysicalOp::PathBuild(PathBuildExec {
511 input,
512 output,
513 node_vars,
514 rel_vars,
515 shortest_path_all,
516 }) => {
517 let upstream = build_streaming(plan, *input, storage, params.clone())?;
518 let ctx = StreamCtx::new(storage, params);
519 Ok(Box::new(PathBuildSource::new(
520 upstream,
521 ctx,
522 *output,
523 node_vars,
524 rel_vars,
525 *shortest_path_all,
526 )))
527 }
528
529 _ => unreachable!("non-streaming op reached streaming branch: {op:?}"),
531 }
532}
533
534fn open_input<'a, S: GraphStorage + 'a>(
538 plan: &'a PhysicalPlan,
539 input: Option<PhysicalNodeId>,
540 storage: &'a S,
541 params: Arc<BTreeMap<String, LoraValue>>,
542) -> ExecResult<Box<dyn RowSource + 'a>> {
543 match input {
544 Some(input) => build_streaming(plan, input, storage, params),
545 None => Ok(Box::new(ArgumentSource::new())),
546 }
547}
548
549fn build_buffered_subtree<'a, S: GraphStorage + 'a>(
556 plan: &'a PhysicalPlan,
557 node_id: PhysicalNodeId,
558 storage: &'a S,
559 params: &Arc<BTreeMap<String, LoraValue>>,
560) -> ExecResult<Box<dyn RowSource + 'a>> {
561 let executor = Executor::new(ExecutionContext {
565 storage,
566 params: (**params).clone(),
567 });
568 let rows = executor.execute_subtree(plan, node_id)?;
569 Ok(Box::new(BufferedRowSource::new(rows)))
570}
571
572pub struct PullExecutor<'a, S: GraphStorage> {
578 storage: &'a S,
579 params: BTreeMap<String, LoraValue>,
580}
581
582impl<'a, S: GraphStorage> PullExecutor<'a, S> {
583 pub fn new(storage: &'a S, params: BTreeMap<String, LoraValue>) -> Self {
584 Self { storage, params }
585 }
586
587 pub fn open_compiled(self, compiled: &'a CompiledQuery) -> ExecResult<Box<dyn RowSource + 'a>>
596 where
597 S: 'a,
598 {
599 clear_eval_error();
600 compiled_to_streaming(compiled, self.storage, self.params)
601 }
602}
603
604pub struct MutablePullExecutor<'a, S: GraphStorageMut> {
609 storage: &'a mut S,
610 params: BTreeMap<String, LoraValue>,
611}
612
613impl<'a, S: GraphStorageMut + GraphStorage> MutablePullExecutor<'a, S> {
614 pub fn new(storage: &'a mut S, params: BTreeMap<String, LoraValue>) -> Self {
615 Self { storage, params }
616 }
617
618 pub fn open_compiled(self, compiled: &'a CompiledQuery) -> ExecResult<Box<dyn RowSource + 'a>>
632 where
633 S: 'a,
634 {
635 if compiled.unions.is_empty() {
636 return open_mutable_plan_cursor(self.storage, &compiled.physical, self.params);
637 }
638
639 MutableUnionSource::open(self.storage, compiled, self.params)
640 .map(|source| Box::new(source) as Box<dyn RowSource + 'a>)
641 }
642}
643
644fn open_mutable_plan_cursor<'a, S: GraphStorageMut + GraphStorage + 'a>(
645 storage: &'a mut S,
646 plan: &'a PhysicalPlan,
647 params: BTreeMap<String, LoraValue>,
648) -> ExecResult<Box<dyn RowSource + 'a>> {
649 if let Some(input) = write_op_input(plan, plan.root) {
650 if subtree_is_fully_streaming(plan, input) {
651 return StreamingWriteCursor::open(storage, plan, plan.root, params)
652 .map(|c| Box::new(c) as Box<dyn RowSource + 'a>);
653 }
654 }
655
656 let mut executor = MutableExecutor::new(MutableExecutionContext { storage, params });
657 let rows = executor.execute_rows(plan)?;
658 Ok(Box::new(BufferedRowSource::new(rows)))
659}
660
661#[derive(Clone, Copy)]
662struct StoragePtr<S> {
663 ptr: *mut S,
664}
665
666impl<S> StoragePtr<S> {
667 fn from_mut(storage: &mut S) -> Self {
668 Self {
669 ptr: storage as *mut S,
670 }
671 }
672
673 unsafe fn as_ref<'a>(&self) -> &'a S {
674 unsafe { &*self.ptr }
675 }
676
677 unsafe fn as_mut<'a>(&self) -> &'a mut S {
678 unsafe { &mut *self.ptr }
679 }
680}
681
682pub struct MutableUnionSource<'a, S: GraphStorageMut + GraphStorage + 'a> {
686 storage_ptr: StoragePtr<S>,
687 compiled: &'a CompiledQuery,
688 params: BTreeMap<String, LoraValue>,
689 branch_idx: usize,
690 current: Option<Box<dyn RowSource + 'a>>,
691 needs_dedup: bool,
692 seen: BTreeSet<Vec<(String, GroupValueKey)>>,
693 _phantom: std::marker::PhantomData<&'a mut S>,
694}
695
696impl<'a, S: GraphStorageMut + GraphStorage + 'a> MutableUnionSource<'a, S> {
697 fn open(
698 storage: &'a mut S,
699 compiled: &'a CompiledQuery,
700 params: BTreeMap<String, LoraValue>,
701 ) -> ExecResult<Self> {
702 let needs_dedup = compiled.unions.iter().any(|branch| !branch.all);
703 Ok(Self {
704 storage_ptr: StoragePtr::from_mut(storage),
705 compiled,
706 params,
707 branch_idx: 0,
708 current: None,
709 needs_dedup,
710 seen: BTreeSet::new(),
711 _phantom: std::marker::PhantomData,
712 })
713 }
714
715 fn branch_count(&self) -> usize {
716 self.compiled.unions.len() + 1
717 }
718
719 fn branch_plan(&self, idx: usize) -> &'a PhysicalPlan {
720 if idx == 0 {
721 &self.compiled.physical
722 } else {
723 &self.compiled.unions[idx - 1].physical
724 }
725 }
726
727 fn open_branch(&mut self, idx: usize) -> ExecResult<Box<dyn RowSource + 'a>> {
728 let plan = self.branch_plan(idx);
729 let storage = unsafe { self.storage_ptr.as_mut() };
734 open_mutable_plan_cursor(storage, plan, self.params.clone())
735 }
736}
737
738impl<'a, S: GraphStorageMut + GraphStorage + 'a> RowSource for MutableUnionSource<'a, S> {
739 fn next_row(&mut self) -> ExecResult<Option<Row>> {
740 loop {
741 if self.branch_idx >= self.branch_count() {
742 return Ok(None);
743 }
744
745 if self.current.is_none() {
746 self.current = Some(self.open_branch(self.branch_idx)?);
747 }
748
749 match self
750 .current
751 .as_mut()
752 .expect("current branch initialized above")
753 .next_row()?
754 {
755 Some(row) => {
756 if self.needs_dedup {
757 let key = row
758 .iter_named()
759 .map(|(_, name, val)| {
760 (name.into_owned(), GroupValueKey::from_value(val))
761 })
762 .collect();
763 if !self.seen.insert(key) {
764 continue;
765 }
766 }
767 return Ok(Some(row));
768 }
769 None => {
770 self.current.take();
771 self.branch_idx += 1;
772 }
773 }
774 }
775 }
776}
777
778pub struct StreamingWriteCursor<'a, S: GraphStorageMut + GraphStorage + 'a> {
802 upstream: ManuallyDrop<Box<dyn RowSource + 'a>>,
804 storage_ptr: StoragePtr<S>,
807 plan: &'a PhysicalPlan,
809 write_op_node: PhysicalNodeId,
813 params: BTreeMap<String, LoraValue>,
816 _phantom: std::marker::PhantomData<&'a mut S>,
817}
818
819impl<'a, S: GraphStorageMut + GraphStorage + 'a> StreamingWriteCursor<'a, S> {
820 pub(crate) fn open(
824 storage: &'a mut S,
825 plan: &'a PhysicalPlan,
826 write_op_node: PhysicalNodeId,
827 params: BTreeMap<String, LoraValue>,
828 ) -> ExecResult<Self> {
829 let input = match write_op_input(plan, write_op_node) {
830 Some(i) => i,
831 None => {
832 return Err(ExecutorError::RuntimeError(format!(
833 "StreamingWriteCursor::open called with non-write node {write_op_node:?}"
834 )));
835 }
836 };
837 let storage_ptr = StoragePtr::from_mut(storage);
838
839 let storage_ref: &'a S = unsafe { storage_ptr.as_ref() };
841 let upstream = build_streaming(plan, input, storage_ref, Arc::new(params.clone()))?;
842
843 Ok(Self {
844 upstream: ManuallyDrop::new(upstream),
845 storage_ptr,
846 plan,
847 write_op_node,
848 params,
849 _phantom: std::marker::PhantomData,
850 })
851 }
852}
853
854impl<'a, S: GraphStorageMut + GraphStorage + 'a> RowSource for StreamingWriteCursor<'a, S> {
855 fn next_row(&mut self) -> ExecResult<Option<Row>> {
856 let mut row = match self.upstream.next_row()? {
857 Some(r) => r,
858 None => return Ok(None),
859 };
860
861 let storage_mut: &mut S = unsafe { self.storage_ptr.as_mut() };
866 let mut exec = MutableExecutor::new(MutableExecutionContext {
867 storage: storage_mut,
868 params: self.params.clone(),
869 });
870 let op = &self.plan.nodes[self.write_op_node];
871 exec.apply_write_op(op, &mut row)?;
872 let row = exec.hydrate_row(row);
873 Ok(Some(row))
874 }
875}
876
877impl<'a, S: GraphStorageMut + GraphStorage + 'a> Drop for StreamingWriteCursor<'a, S> {
878 fn drop(&mut self) {
879 unsafe {
883 ManuallyDrop::drop(&mut self.upstream);
884 }
885 }
886}
887
888pub fn collect_compiled<'a, S: GraphStorage + 'a>(
891 storage: &'a S,
892 params: BTreeMap<String, LoraValue>,
893 compiled: &'a CompiledQuery,
894) -> ExecResult<Vec<Row>> {
895 let mut cursor = PullExecutor::new(storage, params).open_compiled(compiled)?;
896 drain(cursor.as_mut())
897}
898
899#[derive(Debug, Clone, Copy, PartialEq, Eq)]
906pub enum StreamShape {
907 ReadOnly,
910 Mutating,
914}
915
916impl StreamShape {
917 pub fn is_mutating(self) -> bool {
918 matches!(self, StreamShape::Mutating)
919 }
920}
921
922fn plan_is_mutating(plan: &PhysicalPlan) -> bool {
923 plan.nodes.iter().any(|op| {
924 matches!(
925 op,
926 PhysicalOp::Create(_)
927 | PhysicalOp::Merge(_)
928 | PhysicalOp::Delete(_)
929 | PhysicalOp::Set(_)
930 | PhysicalOp::Remove(_)
931 )
932 })
933}
934
935pub fn classify_stream(compiled: &CompiledQuery) -> StreamShape {
939 if plan_is_mutating(&compiled.physical)
940 || compiled
941 .unions
942 .iter()
943 .any(|b| plan_is_mutating(&b.physical))
944 {
945 StreamShape::Mutating
946 } else {
947 StreamShape::ReadOnly
948 }
949}
950
951pub fn plan_result_columns(plan: &PhysicalPlan) -> Vec<String> {
965 plan_columns_at(plan, plan.root).unwrap_or_default()
966}
967
968fn plan_columns_at(plan: &PhysicalPlan, node: PhysicalNodeId) -> Option<Vec<String>> {
969 match &plan.nodes[node] {
970 PhysicalOp::Projection(p) => Some(p.items.iter().map(|i| i.name.clone()).collect()),
971 PhysicalOp::HashAggregation(p) => Some(
972 p.group_by
973 .iter()
974 .chain(p.aggregates.iter())
975 .map(|i| i.name.clone())
976 .collect(),
977 ),
978 PhysicalOp::Limit(p) => plan_columns_at(plan, p.input),
979 PhysicalOp::Sort(p) => plan_columns_at(plan, p.input),
980 PhysicalOp::PathBuild(p) => plan_columns_at(plan, p.input),
981 PhysicalOp::OptionalMatch(p) => plan_columns_at(plan, p.input),
982 PhysicalOp::Filter(p) => plan_columns_at(plan, p.input),
983 PhysicalOp::Unwind(p) => plan_columns_at(plan, p.input),
984 PhysicalOp::Create(p) => plan_columns_at(plan, p.input),
985 PhysicalOp::Merge(p) => plan_columns_at(plan, p.input),
986 PhysicalOp::Delete(p) => plan_columns_at(plan, p.input),
987 PhysicalOp::Set(p) => plan_columns_at(plan, p.input),
988 PhysicalOp::Remove(p) => plan_columns_at(plan, p.input),
989 PhysicalOp::Argument(_)
990 | PhysicalOp::NodeScan(_)
991 | PhysicalOp::NodeByLabelScan(_)
992 | PhysicalOp::NodeByPropertyScan(_)
993 | PhysicalOp::Expand(_) => None,
994 }
995}
996
997pub fn compiled_result_columns(compiled: &CompiledQuery) -> Vec<String> {
1000 plan_result_columns(&compiled.physical)
1001}