use std::collections::{BTreeMap, BTreeSet};
use std::mem::ManuallyDrop;
use std::sync::Arc;
use lora_compiler::physical::{
ExpandExec, FilterExec, HashAggregationExec, LimitExec, NodeByLabelScanExec,
NodeByPropertyScanExec, NodeScanExec, OptionalMatchExec, PathBuildExec, PhysicalNodeId,
PhysicalOp, PhysicalPlan, ProjectionExec, SortExec, UnwindExec,
};
use lora_compiler::CompiledQuery;
use lora_store::{GraphStorage, GraphStorageMut};
use crate::errors::{ExecResult, ExecutorError};
use crate::eval::{clear_eval_error, eval_expr, EvalContext};
use crate::executor::{
hydrate_node_record, hydrate_relationship_record, ExecutionContext, Executor, GroupValueKey,
MutableExecutionContext, MutableExecutor,
};
use crate::value::{LoraValue, Row};
use super::aggregate::HashAggregationSource;
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::{NodeByLabelScanSource, NodeByPropertyScanSource, NodeScanSource};
use super::sort::{LimitSource, SortSource};
use super::union::UnionSource;
pub trait RowSource {
fn next_row(&mut self) -> ExecResult<Option<Row>>;
}
pub fn drain<S: RowSource + ?Sized>(source: &mut S) -> ExecResult<Vec<Row>> {
let mut out = Vec::new();
while let Some(row) = source.next_row()? {
out.push(row);
}
Ok(out)
}
#[derive(Clone)]
pub(super) struct StreamCtx<'a, S: GraphStorage> {
pub storage: &'a S,
pub params: Arc<BTreeMap<String, LoraValue>>,
}
impl<'a, S: GraphStorage> StreamCtx<'a, S> {
pub(super) fn new(storage: &'a S, params: Arc<BTreeMap<String, LoraValue>>) -> Self {
Self { storage, params }
}
pub(super) fn eval_ctx<'b>(&'b self) -> EvalContext<'b, S> {
EvalContext {
storage: self.storage,
params: &self.params,
}
}
}
pub struct BufferedRowSource {
iter: std::vec::IntoIter<Row>,
}
impl BufferedRowSource {
pub fn new(rows: Vec<Row>) -> Self {
Self {
iter: rows.into_iter(),
}
}
}
impl RowSource for BufferedRowSource {
fn next_row(&mut self) -> ExecResult<Option<Row>> {
Ok(self.iter.next())
}
}
pub struct ArgumentSource {
yielded: bool,
}
impl ArgumentSource {
pub fn new() -> Self {
Self { yielded: false }
}
}
impl Default for ArgumentSource {
fn default() -> Self {
Self::new()
}
}
impl RowSource for ArgumentSource {
fn next_row(&mut self) -> ExecResult<Option<Row>> {
if self.yielded {
Ok(None)
} else {
self.yielded = true;
Ok(Some(Row::new()))
}
}
}
pub struct HydratingSource<'a, S: GraphStorage> {
upstream: Box<dyn RowSource + 'a>,
storage: &'a S,
}
impl<'a, S: GraphStorage> HydratingSource<'a, S> {
pub(super) fn new(upstream: Box<dyn RowSource + 'a>, storage: &'a S) -> Self {
Self { upstream, storage }
}
}
impl<'a, S: GraphStorage> RowSource for HydratingSource<'a, S> {
fn next_row(&mut self) -> ExecResult<Option<Row>> {
match self.upstream.next_row()? {
None => Ok(None),
Some(row) => {
let mut out = Row::new();
for (var, name, value) in row.into_iter_named() {
out.insert_named(var, name, hydrate_value(value, self.storage));
}
Ok(Some(out))
}
}
}
}
pub(super) fn hydrate_value<S: GraphStorage>(value: LoraValue, storage: &S) -> LoraValue {
match value {
LoraValue::Node(id) => storage
.with_node(id, hydrate_node_record)
.unwrap_or(LoraValue::Null),
LoraValue::Relationship(id) => storage
.with_relationship(id, hydrate_relationship_record)
.unwrap_or(LoraValue::Null),
LoraValue::List(values) => LoraValue::List(
values
.into_iter()
.map(|v| hydrate_value(v, storage))
.collect(),
),
LoraValue::Map(map) => LoraValue::Map(
map.into_iter()
.map(|(k, v)| (k, hydrate_value(v, storage)))
.collect(),
),
other => other,
}
}
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)?;
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(),
)?;
branches.push(Box::new(HydratingSource::new(head_inner, storage)));
let mut needs_dedup = false;
for branch in &compiled.unions {
let inner = build_streaming(
&branch.physical,
branch.physical.root,
storage,
params.clone(),
)?;
branches.push(Box::new(HydratingSource::new(inner, storage)));
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::Filter(_)
| PhysicalOp::Unwind(_)
| PhysicalOp::Limit(_)
| PhysicalOp::Sort(_)
| PhysicalOp::HashAggregation(_)
| PhysicalOp::OptionalMatch(_)
| PhysicalOp::PathBuild(_)
| PhysicalOp::Projection(_) => true,
PhysicalOp::Expand(_) => true,
_ => false,
}
}
pub(super) 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::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::PathBuild(o) => Some(o.input),
_ => return false,
};
match child {
None => true,
Some(c) => subtree_is_fully_streaming(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>> {
let op = &plan.nodes[node_id];
if !is_streaming_op(op) {
return build_buffered_subtree(plan, node_id, storage, ¶ms);
}
match op {
PhysicalOp::Argument(_) => Ok(Box::new(ArgumentSource::new())),
PhysicalOp::NodeScan(NodeScanExec { input, var }) => {
let upstream = open_input(plan, *input, storage, params.clone())?;
Ok(Box::new(NodeScanSource::new(upstream, storage, *var)))
}
PhysicalOp::NodeByLabelScan(NodeByLabelScanExec { input, var, labels }) => {
let upstream = open_input(plan, *input, storage, params.clone())?;
Ok(Box::new(NodeByLabelScanSource::new(
upstream, storage, *var, labels,
)))
}
PhysicalOp::NodeByPropertyScan(NodeByPropertyScanExec {
input,
var,
labels,
key,
value,
}) => {
let upstream = open_input(plan, *input, storage, params.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(NodeByPropertyScanSource::new(
upstream, ctx, *var, labels, key, value,
)))
}
PhysicalOp::Expand(ExpandExec {
input,
src,
rel,
dst,
types,
direction,
rel_properties,
range,
}) => {
let upstream = build_streaming(plan, *input, storage, params.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(plan, *input, storage, params.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(plan, *input, storage, params.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(plan, *input, storage, params.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(plan, *input, storage, params.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 }) => {
let upstream = build_streaming(plan, *input, storage, params.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(SortSource::new(upstream, ctx, items)))
}
PhysicalOp::HashAggregation(HashAggregationExec {
input,
group_by,
aggregates,
}) => {
let upstream = build_streaming(plan, *input, storage, params.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(plan, *input, storage, params.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(OptionalMatchSource::new(
upstream, ctx, plan, *inner, new_vars,
)))
}
PhysicalOp::PathBuild(PathBuildExec {
input,
output,
node_vars,
rel_vars,
shortest_path_all,
}) => {
let upstream = build_streaming(plan, *input, storage, params.clone())?;
let ctx = StreamCtx::new(storage, params);
Ok(Box::new(PathBuildSource::new(
upstream,
ctx,
*output,
node_vars,
rel_vars,
*shortest_path_all,
)))
}
_ => unreachable!("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>>,
) -> ExecResult<Box<dyn RowSource + 'a>> {
match input {
Some(input) => build_streaming(plan, input, storage, params),
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::new(ExecutionContext {
storage,
params: (**params).clone(),
});
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 struct MutablePullExecutor<'a, S: GraphStorageMut> {
storage: &'a mut S,
params: BTreeMap<String, LoraValue>,
}
impl<'a, S: GraphStorageMut + GraphStorage> MutablePullExecutor<'a, S> {
pub fn new(storage: &'a mut 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,
{
if compiled.unions.is_empty() {
return open_mutable_plan_cursor(self.storage, &compiled.physical, self.params);
}
MutableUnionSource::open(self.storage, compiled, self.params)
.map(|source| Box::new(source) as Box<dyn RowSource + 'a>)
}
}
fn open_mutable_plan_cursor<'a, S: GraphStorageMut + GraphStorage + 'a>(
storage: &'a mut S,
plan: &'a PhysicalPlan,
params: BTreeMap<String, LoraValue>,
) -> ExecResult<Box<dyn RowSource + 'a>> {
if let Some(input) = write_op_input(plan, plan.root) {
if subtree_is_fully_streaming(plan, input) {
return StreamingWriteCursor::open(storage, plan, plan.root, params)
.map(|c| Box::new(c) as Box<dyn RowSource + 'a>);
}
}
let mut executor = MutableExecutor::new(MutableExecutionContext { storage, params });
let rows = executor.execute_rows(plan)?;
Ok(Box::new(BufferedRowSource::new(rows)))
}
#[derive(Clone, Copy)]
struct StoragePtr<S> {
ptr: *mut S,
}
impl<S> StoragePtr<S> {
fn from_mut(storage: &mut S) -> Self {
Self {
ptr: storage as *mut S,
}
}
unsafe fn as_ref<'a>(&self) -> &'a S {
unsafe { &*self.ptr }
}
unsafe fn as_mut<'a>(&self) -> &'a mut S {
unsafe { &mut *self.ptr }
}
}
pub struct MutableUnionSource<'a, S: GraphStorageMut + GraphStorage + 'a> {
storage_ptr: StoragePtr<S>,
compiled: &'a CompiledQuery,
params: BTreeMap<String, LoraValue>,
branch_idx: usize,
current: Option<Box<dyn RowSource + 'a>>,
needs_dedup: bool,
seen: BTreeSet<Vec<(String, GroupValueKey)>>,
_phantom: std::marker::PhantomData<&'a mut S>,
}
impl<'a, S: GraphStorageMut + GraphStorage + 'a> MutableUnionSource<'a, S> {
fn open(
storage: &'a mut S,
compiled: &'a CompiledQuery,
params: BTreeMap<String, LoraValue>,
) -> ExecResult<Self> {
let needs_dedup = compiled.unions.iter().any(|branch| !branch.all);
Ok(Self {
storage_ptr: StoragePtr::from_mut(storage),
compiled,
params,
branch_idx: 0,
current: None,
needs_dedup,
seen: BTreeSet::new(),
_phantom: std::marker::PhantomData,
})
}
fn branch_count(&self) -> usize {
self.compiled.unions.len() + 1
}
fn branch_plan(&self, idx: usize) -> &'a PhysicalPlan {
if idx == 0 {
&self.compiled.physical
} else {
&self.compiled.unions[idx - 1].physical
}
}
fn open_branch(&mut self, idx: usize) -> ExecResult<Box<dyn RowSource + 'a>> {
let plan = self.branch_plan(idx);
let storage = unsafe { self.storage_ptr.as_mut() };
open_mutable_plan_cursor(storage, plan, self.params.clone())
}
}
impl<'a, S: GraphStorageMut + GraphStorage + 'a> RowSource for MutableUnionSource<'a, S> {
fn next_row(&mut self) -> ExecResult<Option<Row>> {
loop {
if self.branch_idx >= self.branch_count() {
return Ok(None);
}
if self.current.is_none() {
self.current = Some(self.open_branch(self.branch_idx)?);
}
match self
.current
.as_mut()
.expect("current branch initialized above")
.next_row()?
{
Some(row) => {
if self.needs_dedup {
let key = row
.iter_named()
.map(|(_, name, val)| {
(name.into_owned(), GroupValueKey::from_value(val))
})
.collect();
if !self.seen.insert(key) {
continue;
}
}
return Ok(Some(row));
}
None => {
self.current.take();
self.branch_idx += 1;
}
}
}
}
}
pub struct StreamingWriteCursor<'a, S: GraphStorageMut + GraphStorage + 'a> {
upstream: ManuallyDrop<Box<dyn RowSource + 'a>>,
storage_ptr: StoragePtr<S>,
plan: &'a PhysicalPlan,
write_op_node: PhysicalNodeId,
params: BTreeMap<String, LoraValue>,
_phantom: std::marker::PhantomData<&'a mut S>,
}
impl<'a, S: GraphStorageMut + GraphStorage + 'a> StreamingWriteCursor<'a, S> {
pub(crate) fn open(
storage: &'a mut S,
plan: &'a PhysicalPlan,
write_op_node: PhysicalNodeId,
params: BTreeMap<String, LoraValue>,
) -> ExecResult<Self> {
let input = match write_op_input(plan, write_op_node) {
Some(i) => i,
None => {
return Err(ExecutorError::RuntimeError(format!(
"StreamingWriteCursor::open called with non-write node {write_op_node:?}"
)));
}
};
let storage_ptr = StoragePtr::from_mut(storage);
let storage_ref: &'a S = unsafe { storage_ptr.as_ref() };
let upstream = build_streaming(plan, input, storage_ref, Arc::new(params.clone()))?;
Ok(Self {
upstream: ManuallyDrop::new(upstream),
storage_ptr,
plan,
write_op_node,
params,
_phantom: std::marker::PhantomData,
})
}
}
impl<'a, S: GraphStorageMut + GraphStorage + 'a> RowSource for StreamingWriteCursor<'a, S> {
fn next_row(&mut self) -> ExecResult<Option<Row>> {
let mut row = match self.upstream.next_row()? {
Some(r) => r,
None => return Ok(None),
};
let storage_mut: &mut S = unsafe { self.storage_ptr.as_mut() };
let mut exec = MutableExecutor::new(MutableExecutionContext {
storage: storage_mut,
params: self.params.clone(),
});
let op = &self.plan.nodes[self.write_op_node];
exec.apply_write_op(op, &mut row)?;
let row = exec.hydrate_row(row);
Ok(Some(row))
}
}
impl<'a, S: GraphStorageMut + GraphStorage + 'a> Drop for StreamingWriteCursor<'a, S> {
fn drop(&mut self) {
unsafe {
ManuallyDrop::drop(&mut self.upstream);
}
}
}
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())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StreamShape {
ReadOnly,
Mutating,
}
impl StreamShape {
pub fn is_mutating(self) -> bool {
matches!(self, StreamShape::Mutating)
}
}
fn plan_is_mutating(plan: &PhysicalPlan) -> bool {
plan.nodes.iter().any(|op| {
matches!(
op,
PhysicalOp::Create(_)
| PhysicalOp::Merge(_)
| PhysicalOp::Delete(_)
| PhysicalOp::Set(_)
| PhysicalOp::Remove(_)
)
})
}
pub fn classify_stream(compiled: &CompiledQuery) -> StreamShape {
if plan_is_mutating(&compiled.physical)
|| compiled
.unions
.iter()
.any(|b| plan_is_mutating(&b.physical))
{
StreamShape::Mutating
} else {
StreamShape::ReadOnly
}
}
pub fn plan_result_columns(plan: &PhysicalPlan) -> Vec<String> {
plan_columns_at(plan, plan.root).unwrap_or_default()
}
fn plan_columns_at(plan: &PhysicalPlan, node: PhysicalNodeId) -> Option<Vec<String>> {
match &plan.nodes[node] {
PhysicalOp::Projection(p) => Some(p.items.iter().map(|i| i.name.clone()).collect()),
PhysicalOp::HashAggregation(p) => Some(
p.group_by
.iter()
.chain(p.aggregates.iter())
.map(|i| i.name.clone())
.collect(),
),
PhysicalOp::Limit(p) => plan_columns_at(plan, p.input),
PhysicalOp::Sort(p) => plan_columns_at(plan, p.input),
PhysicalOp::PathBuild(p) => plan_columns_at(plan, p.input),
PhysicalOp::OptionalMatch(p) => plan_columns_at(plan, p.input),
PhysicalOp::Filter(p) => plan_columns_at(plan, p.input),
PhysicalOp::Unwind(p) => plan_columns_at(plan, p.input),
PhysicalOp::Create(p) => plan_columns_at(plan, p.input),
PhysicalOp::Merge(p) => plan_columns_at(plan, p.input),
PhysicalOp::Delete(p) => plan_columns_at(plan, p.input),
PhysicalOp::Set(p) => plan_columns_at(plan, p.input),
PhysicalOp::Remove(p) => plan_columns_at(plan, p.input),
PhysicalOp::Argument(_)
| PhysicalOp::NodeScan(_)
| PhysicalOp::NodeByLabelScan(_)
| PhysicalOp::NodeByPropertyScan(_)
| PhysicalOp::Expand(_) => None,
}
}
pub fn compiled_result_columns(compiled: &CompiledQuery) -> Vec<String> {
plan_result_columns(&compiled.physical)
}