use std::collections::BTreeMap;
use std::sync::Arc;
use lora_compiler::physical::{
CallSubqueryExec, ExpandExec, FilterExec, HashAggregationExec, LimitExec, NodeByLabelScanExec,
NodeByPointScanExec, NodeByPropertyRangeScanExec, NodeByPropertyScanExec, NodeByTextScanExec,
NodeScanExec, OptionalMatchExec, PathBuildExec, PhysicalNodeId, PhysicalOp, PhysicalPlan,
ProjectionExec, RelByPointScanExec, RelByPropertyRangeScanExec, RelByTextScanExec, SortExec,
UnwindExec,
};
use lora_compiler::CompiledQuery;
use lora_store::GraphStorage;
use crate::errors::{ExecResult, ExecutorError};
use crate::eval::{clear_eval_error, eval_expr};
use crate::executor::{plan_may_need_hydration, ExecutionContext, Executor};
use crate::profile::wrap_metered;
use crate::value::{LoraValue, Row};
use super::aggregate::HashAggregationSource;
use super::call_subquery::CallSubquerySource;
use super::expand::{ExpandSource, VariableLengthExpandSource};
use super::filter::FilterSource;
use super::optional::OptionalMatchSource;
use super::path::PathBuildSource;
use super::projection::{DistinctSource, ProjectionSource, UnwindSource};
use super::scan::{
BufferedIndexScanSource, NodeByLabelScanSource, NodeByPropertyScanSource, NodeScanSource,
};
use super::sort::{LimitSource, SortSource};
use super::union::UnionSource;
use super::{drain, ArgumentSource, BufferedRowSource, HydratingSource, RowSource, StreamCtx};
pub(super) fn compiled_to_streaming<'a, S: GraphStorage + 'a>(
compiled: &'a CompiledQuery,
storage: &'a S,
params: BTreeMap<String, LoraValue>,
) -> ExecResult<Box<dyn RowSource + 'a>> {
let params = Arc::new(params);
if compiled.unions.is_empty() {
let plan = &compiled.physical;
let inner = build_streaming(plan, plan.root, storage, params)?;
if !plan_may_need_hydration(plan) {
return Ok(inner);
}
return Ok(Box::new(HydratingSource::new(inner, storage)));
}
let mut branches: Vec<Box<dyn RowSource + 'a>> = Vec::with_capacity(compiled.unions.len() + 1);
let head_inner = build_streaming(
&compiled.physical,
compiled.physical.root,
storage,
params.clone(),
)?;
if plan_may_need_hydration(&compiled.physical) {
branches.push(Box::new(HydratingSource::new(head_inner, storage)));
} else {
branches.push(head_inner);
}
let mut needs_dedup = false;
for branch in &compiled.unions {
let inner = build_streaming(
&branch.physical,
branch.physical.root,
storage,
params.clone(),
)?;
if plan_may_need_hydration(&branch.physical) {
branches.push(Box::new(HydratingSource::new(inner, storage)));
} else {
branches.push(inner);
}
if !branch.all {
needs_dedup = true;
}
}
Ok(Box::new(UnionSource::new(branches, needs_dedup)))
}
pub(super) fn is_streaming_op(op: &PhysicalOp) -> bool {
match op {
PhysicalOp::Argument(_)
| PhysicalOp::NodeScan(_)
| PhysicalOp::NodeByLabelScan(_)
| PhysicalOp::NodeByPropertyScan(_)
| PhysicalOp::NodeByPropertyRangeScan(_)
| PhysicalOp::NodeByTextScan(_)
| PhysicalOp::NodeByPointScan(_)
| PhysicalOp::RelByPropertyRangeScan(_)
| PhysicalOp::RelByTextScan(_)
| PhysicalOp::RelByPointScan(_)
| PhysicalOp::Filter(_)
| PhysicalOp::Unwind(_)
| PhysicalOp::Limit(_)
| PhysicalOp::Sort(_)
| PhysicalOp::HashAggregation(_)
| PhysicalOp::OptionalMatch(_)
| PhysicalOp::PathBuild(_)
| PhysicalOp::CallSubquery(_)
| PhysicalOp::Projection(_) => true,
PhysicalOp::Expand(_) => true,
_ => false,
}
}
pub(crate) fn write_op_input(
plan: &PhysicalPlan,
node_id: PhysicalNodeId,
) -> Option<PhysicalNodeId> {
match &plan.nodes[node_id] {
PhysicalOp::Create(o) => Some(o.input),
PhysicalOp::Set(o) => Some(o.input),
PhysicalOp::Delete(o) => Some(o.input),
PhysicalOp::Remove(o) => Some(o.input),
PhysicalOp::Merge(o) => Some(o.input),
_ => None,
}
}
pub(crate) fn subtree_is_fully_streaming(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
let op = &plan.nodes[node_id];
if !is_streaming_op(op) {
return false;
}
let child = match op {
PhysicalOp::Argument(_) => return true,
PhysicalOp::NodeScan(o) => o.input,
PhysicalOp::NodeByLabelScan(o) => o.input,
PhysicalOp::NodeByPropertyScan(o) => o.input,
PhysicalOp::NodeByPropertyRangeScan(o) => o.input,
PhysicalOp::NodeByTextScan(o) => o.input,
PhysicalOp::NodeByPointScan(o) => o.input,
PhysicalOp::RelByPropertyRangeScan(o) => o.input,
PhysicalOp::RelByTextScan(o) => o.input,
PhysicalOp::RelByPointScan(o) => o.input,
PhysicalOp::Filter(o) => Some(o.input),
PhysicalOp::Unwind(o) => Some(o.input),
PhysicalOp::Limit(o) => Some(o.input),
PhysicalOp::Expand(o) => Some(o.input),
PhysicalOp::Projection(o) => Some(o.input),
PhysicalOp::Sort(o) => Some(o.input),
PhysicalOp::HashAggregation(o) => Some(o.input),
PhysicalOp::OptionalMatch(o) => Some(o.input),
PhysicalOp::CallSubquery(o) if subtree_has_write(plan, o.inner) => return false,
PhysicalOp::CallSubquery(o) => Some(o.input),
PhysicalOp::PathBuild(o) => Some(o.input),
_ => return false,
};
match child {
None => true,
Some(c) => subtree_is_fully_streaming(plan, c),
}
}
pub(crate) fn subtree_has_write(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
let op = &plan.nodes[node_id];
let (first, second) = match op {
PhysicalOp::Create(_)
| PhysicalOp::Merge(_)
| PhysicalOp::Delete(_)
| PhysicalOp::Set(_)
| PhysicalOp::Remove(_)
| PhysicalOp::Foreach(_) => return true,
PhysicalOp::Argument(_) => (None, None),
PhysicalOp::NodeScan(o) => (o.input, None),
PhysicalOp::NodeByLabelScan(o) => (o.input, None),
PhysicalOp::NodeByPropertyScan(o) => (o.input, None),
PhysicalOp::NodeByPropertyRangeScan(o) => (o.input, None),
PhysicalOp::NodeByTextScan(o) => (o.input, None),
PhysicalOp::NodeByPointScan(o) => (o.input, None),
PhysicalOp::RelByPropertyRangeScan(o) => (o.input, None),
PhysicalOp::RelByTextScan(o) => (o.input, None),
PhysicalOp::RelByPointScan(o) => (o.input, None),
PhysicalOp::Expand(o) => (Some(o.input), None),
PhysicalOp::Filter(o) => (Some(o.input), None),
PhysicalOp::Projection(o) => (Some(o.input), None),
PhysicalOp::Unwind(o) => (Some(o.input), None),
PhysicalOp::HashAggregation(o) => (Some(o.input), None),
PhysicalOp::Sort(o) => (Some(o.input), None),
PhysicalOp::Limit(o) => (Some(o.input), None),
PhysicalOp::PathBuild(o) => (Some(o.input), None),
PhysicalOp::OptionalMatch(o) => (Some(o.input), Some(o.inner)),
PhysicalOp::CallSubquery(o) => (Some(o.input), Some(o.inner)),
};
first.is_some_and(|c| subtree_has_write(plan, c))
|| second.is_some_and(|c| subtree_has_write(plan, c))
}
pub(crate) fn build_streaming<'a, S: GraphStorage + 'a>(
plan: &'a PhysicalPlan,
node_id: PhysicalNodeId,
storage: &'a S,
params: Arc<BTreeMap<String, LoraValue>>,
) -> ExecResult<Box<dyn RowSource + 'a>> {
build_streaming_dispatch(plan, node_id, storage, params, None)
}
pub(crate) fn build_streaming_seeded<'a, S: GraphStorage + 'a>(
plan: &'a PhysicalPlan,
node_id: PhysicalNodeId,
storage: &'a S,
params: Arc<BTreeMap<String, LoraValue>>,
seed: Row,
) -> ExecResult<Box<dyn RowSource + 'a>> {
build_streaming_dispatch(plan, node_id, storage, params, Some(seed))
}
fn build_streaming_dispatch<'a, S: GraphStorage + 'a>(
plan: &'a PhysicalPlan,
node_id: PhysicalNodeId,
storage: &'a S,
params: Arc<BTreeMap<String, LoraValue>>,
seed: Option<Row>,
) -> ExecResult<Box<dyn RowSource + 'a>> {
build_streaming_inner(plan, node_id, storage, params, seed).map(|src| {
super::source::DeadlineSource::wrap(
wrap_metered(node_id, src),
crate::cancel::active_deadline(),
)
})
}
fn build_streaming_inner<'a, S: GraphStorage + 'a>(
plan: &'a PhysicalPlan,
node_id: PhysicalNodeId,
storage: &'a S,
params: Arc<BTreeMap<String, LoraValue>>,
seed: Option<Row>,
) -> ExecResult<Box<dyn RowSource + 'a>> {
let op = &plan.nodes[node_id];
if !is_streaming_op(op) {
return build_buffered_subtree(plan, node_id, storage, ¶ms);
}
match op {
PhysicalOp::Argument(_) => match seed {
Some(seed_row) => Ok(Box::new(BufferedRowSource::new(vec![seed_row]))),
None => Ok(Box::new(ArgumentSource::new())),
},
PhysicalOp::NodeScan(NodeScanExec { input, var }) => {
let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
Ok(Box::new(NodeScanSource::new(upstream, storage, *var)))
}
PhysicalOp::NodeByLabelScan(NodeByLabelScanExec { input, var, labels }) => {
let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
Ok(Box::new(NodeByLabelScanSource::new(
upstream, storage, *var, labels,
)))
}
PhysicalOp::NodeByPropertyScan(NodeByPropertyScanExec {
input,
var,
labels,
key,
value,
in_list,
}) => {
let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(NodeByPropertyScanSource::new(
upstream, ctx, *var, labels, key, value, *in_list,
)))
}
PhysicalOp::NodeByPropertyRangeScan(op @ NodeByPropertyRangeScanExec { input, .. }) => {
let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
if op.order.is_some() {
return Ok(Box::new(super::scan::OrderedRangeScanSource::new(
upstream, ctx, op,
)));
}
Ok(Box::new(BufferedIndexScanSource::node_range(
upstream, ctx, op,
)))
}
PhysicalOp::NodeByTextScan(op @ NodeByTextScanExec { input, .. }) => {
let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(BufferedIndexScanSource::node_text(
upstream, ctx, op,
)))
}
PhysicalOp::NodeByPointScan(op @ NodeByPointScanExec { input, .. }) => {
let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(BufferedIndexScanSource::node_point(
upstream, ctx, op,
)))
}
PhysicalOp::RelByPropertyRangeScan(op @ RelByPropertyRangeScanExec { input, .. }) => {
let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(BufferedIndexScanSource::rel_range(
upstream, ctx, op,
)))
}
PhysicalOp::RelByTextScan(op @ RelByTextScanExec { input, .. }) => {
let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(BufferedIndexScanSource::rel_text(
upstream, ctx, op,
)))
}
PhysicalOp::RelByPointScan(op @ RelByPointScanExec { input, .. }) => {
let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(BufferedIndexScanSource::rel_point(
upstream, ctx, op,
)))
}
PhysicalOp::Expand(ExpandExec {
input,
src,
rel,
dst,
types,
direction,
rel_properties,
range,
}) => {
let upstream =
build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
match range.as_ref() {
Some(range) => Ok(Box::new(VariableLengthExpandSource::new(
upstream, ctx, *src, *rel, *dst, types, *direction, range,
))),
None => Ok(Box::new(ExpandSource::new(
upstream,
ctx,
*src,
*rel,
*dst,
types,
*direction,
rel_properties.as_ref(),
))),
}
}
PhysicalOp::Filter(FilterExec { input, predicate }) => {
let upstream =
build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(FilterSource::new(upstream, ctx, predicate)))
}
PhysicalOp::Projection(ProjectionExec {
input,
distinct,
items,
include_existing,
}) => {
let upstream =
build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
let proj: Box<dyn RowSource + 'a> = Box::new(ProjectionSource::new(
upstream,
ctx,
items,
*include_existing,
));
if *distinct {
Ok(Box::new(DistinctSource::new(proj)))
} else {
Ok(proj)
}
}
PhysicalOp::Unwind(UnwindExec { input, expr, alias }) => {
let upstream =
build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(UnwindSource::new(upstream, ctx, expr, *alias)))
}
PhysicalOp::Limit(LimitExec { input, skip, limit }) => {
let upstream =
build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
let eval_ctx = ctx.eval_ctx();
let scratch = Row::new();
let skip_n = skip
.as_ref()
.and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
.unwrap_or(0)
.max(0) as usize;
let limit_n = limit
.as_ref()
.and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
.map(|n| n.max(0) as usize);
Ok(Box::new(LimitSource::new(upstream, skip_n, limit_n)))
}
PhysicalOp::Sort(SortExec {
input,
items,
top_k,
}) => {
let upstream =
build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(SortSource::new_with_top_k(
upstream, ctx, items, *top_k,
)))
}
PhysicalOp::HashAggregation(
agg @ HashAggregationExec {
input,
group_by,
aggregates,
},
) => {
if seed.is_none() {
if let Some(rows) =
crate::executor::count_all_scan_aggregation_rows(storage, plan, agg)
{
return Ok(Box::new(BufferedRowSource::new(rows)));
}
}
let upstream =
build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(HashAggregationSource::new(
upstream, ctx, group_by, aggregates,
)))
}
PhysicalOp::OptionalMatch(OptionalMatchExec {
input,
inner,
new_vars,
}) => {
let upstream =
build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(OptionalMatchSource::new(
upstream, ctx, plan, *inner, new_vars,
)))
}
PhysicalOp::CallSubquery(CallSubqueryExec {
input,
inner,
new_vars,
}) => {
let upstream = build_streaming_dispatch(plan, *input, storage, params.clone(), seed)?;
Ok(Box::new(CallSubquerySource::new(
upstream, plan, *inner, storage, params, new_vars,
)))
}
PhysicalOp::PathBuild(PathBuildExec {
input,
output,
node_vars,
rel_vars,
shortest_path_all,
}) => {
let upstream =
build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(PathBuildSource::new(
upstream,
ctx,
*output,
node_vars,
rel_vars,
*shortest_path_all,
)))
}
_ => Err(ExecutorError::RuntimeError(format!(
"non-streaming op reached streaming branch: {op:?}"
))),
}
}
fn open_input<'a, S: GraphStorage + 'a>(
plan: &'a PhysicalPlan,
input: Option<PhysicalNodeId>,
storage: &'a S,
params: Arc<BTreeMap<String, LoraValue>>,
seed: Option<Row>,
) -> ExecResult<Box<dyn RowSource + 'a>> {
match input {
Some(input) => build_streaming_dispatch(plan, input, storage, params, seed),
None => match seed {
Some(seed_row) => Ok(Box::new(BufferedRowSource::new(vec![seed_row]))),
None => Ok(Box::new(ArgumentSource::new())),
},
}
}
fn build_buffered_subtree<'a, S: GraphStorage + 'a>(
plan: &'a PhysicalPlan,
node_id: PhysicalNodeId,
storage: &'a S,
params: &Arc<BTreeMap<String, LoraValue>>,
) -> ExecResult<Box<dyn RowSource + 'a>> {
let executor = Executor::with_deadline(
ExecutionContext {
storage,
params: (**params).clone(),
},
crate::cancel::active_deadline(),
);
let rows = executor.execute_subtree(plan, node_id)?;
Ok(Box::new(BufferedRowSource::new(rows)))
}
pub struct PullExecutor<'a, S: GraphStorage> {
storage: &'a S,
params: BTreeMap<String, LoraValue>,
}
impl<'a, S: GraphStorage> PullExecutor<'a, S> {
pub fn new(storage: &'a S, params: BTreeMap<String, LoraValue>) -> Self {
Self { storage, params }
}
pub fn open_compiled(self, compiled: &'a CompiledQuery) -> ExecResult<Box<dyn RowSource + 'a>>
where
S: 'a,
{
clear_eval_error();
compiled_to_streaming(compiled, self.storage, self.params)
}
}
pub fn collect_compiled<'a, S: GraphStorage + 'a>(
storage: &'a S,
params: BTreeMap<String, LoraValue>,
compiled: &'a CompiledQuery,
) -> ExecResult<Vec<Row>> {
let mut cursor = PullExecutor::new(storage, params).open_compiled(compiled)?;
drain(cursor.as_mut())
}
pub fn collect_compiled_with_deadline<'a, S: GraphStorage + 'a>(
storage: &'a S,
params: BTreeMap<String, LoraValue>,
compiled: &'a CompiledQuery,
deadline: Option<web_time::Instant>,
) -> ExecResult<Vec<Row>> {
let _deadline_scope = crate::cancel::DeadlineScope::enter(deadline);
if let Some(deadline) = deadline {
if crate::cancel::deadline_reached(deadline) {
return Err(ExecutorError::QueryTimeout);
}
}
let mut cursor = PullExecutor::new(storage, params).open_compiled(compiled)?;
drain(cursor.as_mut())
}