1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
// Copyright (c) 2026 Austin Han <austinhan1024@gmail.com>
//
// This file is part of RocksGraph.
//
// RocksGraph is free software: you can redistribute it and/or modify
// it under the terms of the GNU General Public License as published by
// the Free Software Foundation, either version 2 of the License, or
// (at your option) any later version.
//
// RocksGraph is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU General Public License for more details.
//
// You should have received a copy of the GNU General Public License
// along with RocksGraph. If not, see <https://www.gnu.org/licenses/>.
use crate::types::{PIPELINE_PRODUCE_SIZE, SMALL_VECTOR_LENGTH};
use std::rc::Rc;
use smallvec::{smallvec, SmallVec};
use crate::engine::volcano::steps::traits::ExplainNode;
use crate::{
engine::{
context::GraphCtx,
traverser::Traverser,
volcano::{
builder::PhysicalPlan,
steps::traits::{CoreStep, StepRef},
},
},
types::error::StoreError,
};
/// A physical step that implements the `coalesce` logical step.
#[derive(Debug)]
pub struct CoalesceStep {
// ── Upstream link ──
upstream: Option<StepRef>,
// ── Static/Fixed configuration ──
/// The physical plans/branches to evaluate in coalesce.
physical_plans: SmallVec<[PhysicalPlan; SMALL_VECTOR_LENGTH]>,
// ── Dynamic/Runtime execution state ──
/// The parent traverser currently being evaluated.
current_input: Option<Rc<Traverser>>,
/// The index of the current branch plan being evaluated for the active input.
current_plan_idx: usize,
/// The index of the branch plan that successfully yielded results (if any).
winning_plan_idx: Option<usize>,
}
impl CoalesceStep {
/// Creates a new `CoalesceStep` with the given physical sub-plans.
pub fn new(physical_plans: SmallVec<[PhysicalPlan; SMALL_VECTOR_LENGTH]>) -> Self {
Self { upstream: None, physical_plans, current_input: None, current_plan_idx: 0, winning_plan_idx: None }
}
}
impl CoreStep for CoalesceStep {
fn add_upper(&mut self, upstream: StepRef) {
// Sets the upstream step for this coalesce step.
self.upstream = Some(upstream);
}
fn produce(
&mut self,
ctx: &mut dyn GraphCtx,
) -> Result<Option<SmallVec<[Rc<Traverser>; PIPELINE_PRODUCE_SIZE]>>, StoreError> {
// Produces traversers from the first sub-plan that yields results.
loop {
// If we found a winning branch, keep draining it
if let Some(winning_idx) = self.winning_plan_idx {
if let Some(res) = self.physical_plans[winning_idx].next(ctx)? {
return Ok(Some(smallvec![res]));
}
self.current_input = None;
self.winning_plan_idx = None;
}
// Fetch next input from upstream when current is exhausted
if self.current_input.is_none() {
let Some(upstream) = self.upstream.as_ref() else { return Ok(None) };
let Some(t) = upstream.next(ctx)? else { return Ok(None) };
self.current_input = Some(Rc::clone(&t));
self.current_plan_idx = 0;
if let Some(p) = self.physical_plans.first() {
p.reset();
p.inject(smallvec![Rc::clone(&t)]);
}
}
// All branches exhausted for this input — move to next input traverser
if self.current_plan_idx >= self.physical_plans.len() {
self.current_input = None;
continue;
}
// Try the current branch
if let Some(res) = self.physical_plans[self.current_plan_idx].next(ctx)? {
self.winning_plan_idx = Some(self.current_plan_idx);
return Ok(Some(smallvec![res]));
}
// Branch yielded nothing — advance to next branch
self.current_plan_idx += 1;
if self.current_plan_idx < self.physical_plans.len() {
let t = Rc::clone(self.current_input.as_ref().unwrap());
self.physical_plans[self.current_plan_idx].reset();
self.physical_plans[self.current_plan_idx].inject(smallvec![t]);
}
}
}
fn reset(&mut self) {
// Resets the state of this step and all its upstream and sub-plans.
if let Some(up) = &self.upstream {
up.reset();
}
for p in &self.physical_plans {
p.reset();
}
self.current_input = None;
self.current_plan_idx = 0;
self.winning_plan_idx = None;
}
fn upper(&self) -> Option<StepRef> {
// Returns a clone of the upstream step reference.
self.upstream.clone()
}
fn explain(&self) -> ExplainNode {
let children =
self.physical_plans.iter().enumerate().map(|(i, plan)| (format!("branch {}", i), plan.explain())).collect();
ExplainNode::new("CoalesceStep").with_children(children)
}
}