1use std::collections::{BTreeMap, BTreeSet};
2use std::mem::ManuallyDrop;
3use std::sync::Arc;
4
5use lora_compiler::physical::{PhysicalNodeId, PhysicalPlan};
6use lora_compiler::CompiledQuery;
7use lora_store::{GraphStorage, GraphStorageMut};
8
9use crate::errors::{ExecResult, ExecutorError};
10use crate::executor::{GroupValueKey, MutableExecutionContext, MutableExecutor};
11use crate::value::{LoraValue, Row};
12
13use super::traits::write_op_input;
14use super::{build_streaming, subtree_is_fully_streaming, BufferedRowSource, RowSource};
15
16pub struct MutablePullExecutor<'a, S: GraphStorageMut> {
21 storage: &'a mut S,
22 params: BTreeMap<String, LoraValue>,
23}
24
25impl<'a, S: GraphStorageMut + GraphStorage> MutablePullExecutor<'a, S> {
26 pub fn new(storage: &'a mut S, params: BTreeMap<String, LoraValue>) -> Self {
27 Self { storage, params }
28 }
29
30 pub fn open_compiled(self, compiled: &'a CompiledQuery) -> ExecResult<Box<dyn RowSource + 'a>>
44 where
45 S: 'a,
46 {
47 if compiled.unions.is_empty() {
48 return open_mutable_plan_cursor(self.storage, &compiled.physical, self.params);
49 }
50
51 MutableUnionSource::open(self.storage, compiled, self.params)
52 .map(|source| Box::new(source) as Box<dyn RowSource + 'a>)
53 }
54}
55
56fn open_mutable_plan_cursor<'a, S: GraphStorageMut + GraphStorage + 'a>(
57 storage: &'a mut S,
58 plan: &'a PhysicalPlan,
59 params: BTreeMap<String, LoraValue>,
60) -> ExecResult<Box<dyn RowSource + 'a>> {
61 if let Some(input) = write_op_input(plan, plan.root) {
62 if subtree_is_fully_streaming(plan, input) {
63 let cursor = StreamingWriteCursor::open(storage, plan, plan.root, params)?;
64 let cursor: Box<dyn RowSource + 'a> = Box::new(cursor);
65 if crate::executor::plan_ends_in_write(plan) {
66 return Ok(Box::new(DrainSilently {
67 inner: Some(cursor),
68 }));
69 }
70 return Ok(cursor);
71 }
72 }
73
74 let mut executor = MutableExecutor::new(MutableExecutionContext { storage, params });
75 let rows = executor.execute_rows(plan)?;
76 Ok(Box::new(BufferedRowSource::new(rows)))
77}
78
79#[derive(Clone, Copy)]
80struct StoragePtr<S> {
81 ptr: *mut S,
82}
83
84impl<S> StoragePtr<S> {
85 fn from_mut(storage: &mut S) -> Self {
86 Self {
87 ptr: storage as *mut S,
88 }
89 }
90
91 unsafe fn as_ref<'a>(&self) -> &'a S {
92 unsafe { &*self.ptr }
93 }
94
95 unsafe fn as_mut<'a>(&self) -> &'a mut S {
96 unsafe { &mut *self.ptr }
97 }
98}
99
100pub struct MutableUnionSource<'a, S: GraphStorageMut + GraphStorage + 'a> {
104 storage_ptr: StoragePtr<S>,
105 compiled: &'a CompiledQuery,
106 params: BTreeMap<String, LoraValue>,
107 branch_idx: usize,
108 current: Option<Box<dyn RowSource + 'a>>,
109 needs_dedup: bool,
110 seen: BTreeSet<Vec<(String, GroupValueKey)>>,
111 _phantom: std::marker::PhantomData<&'a mut S>,
112}
113
114impl<'a, S: GraphStorageMut + GraphStorage + 'a> MutableUnionSource<'a, S> {
115 fn open(
116 storage: &'a mut S,
117 compiled: &'a CompiledQuery,
118 params: BTreeMap<String, LoraValue>,
119 ) -> ExecResult<Self> {
120 let needs_dedup = compiled.unions.iter().any(|branch| !branch.all);
121 Ok(Self {
122 storage_ptr: StoragePtr::from_mut(storage),
123 compiled,
124 params,
125 branch_idx: 0,
126 current: None,
127 needs_dedup,
128 seen: BTreeSet::new(),
129 _phantom: std::marker::PhantomData,
130 })
131 }
132
133 fn branch_count(&self) -> usize {
134 self.compiled.unions.len() + 1
135 }
136
137 fn branch_plan(&self, idx: usize) -> &'a PhysicalPlan {
138 if idx == 0 {
139 &self.compiled.physical
140 } else {
141 &self.compiled.unions[idx - 1].physical
142 }
143 }
144
145 fn open_branch(&mut self, idx: usize) -> ExecResult<Box<dyn RowSource + 'a>> {
146 let plan = self.branch_plan(idx);
147 let storage = unsafe { self.storage_ptr.as_mut() };
152 open_mutable_plan_cursor(storage, plan, self.params.clone())
153 }
154}
155
156impl<'a, S: GraphStorageMut + GraphStorage + 'a> RowSource for MutableUnionSource<'a, S> {
157 fn next_row(&mut self) -> ExecResult<Option<Row>> {
158 loop {
159 if self.branch_idx >= self.branch_count() {
160 return Ok(None);
161 }
162
163 if self.current.is_none() {
164 self.current = Some(self.open_branch(self.branch_idx)?);
165 }
166
167 let Some(current) = self.current.as_mut() else {
168 return Err(ExecutorError::RuntimeError(
169 "mutable UNION cursor lost its current branch".into(),
170 ));
171 };
172
173 match current.next_row()? {
174 Some(row) => {
175 if self.needs_dedup {
176 let key = row
177 .iter_named()
178 .map(|(_, name, val)| {
179 (name.into_owned(), GroupValueKey::from_value(val))
180 })
181 .collect();
182 if !self.seen.insert(key) {
183 continue;
184 }
185 }
186 return Ok(Some(row));
187 }
188 None => {
189 self.current.take();
190 self.branch_idx += 1;
191 }
192 }
193 }
194 }
195}
196
197pub struct StreamingWriteCursor<'a, S: GraphStorageMut + GraphStorage + 'a> {
221 upstream: ManuallyDrop<Box<dyn RowSource + 'a>>,
223 storage_ptr: StoragePtr<S>,
226 plan: &'a PhysicalPlan,
228 write_op_node: PhysicalNodeId,
232 params: BTreeMap<String, LoraValue>,
235 _phantom: std::marker::PhantomData<&'a mut S>,
236}
237
238impl<'a, S: GraphStorageMut + GraphStorage + 'a> StreamingWriteCursor<'a, S> {
239 pub(crate) fn open(
243 storage: &'a mut S,
244 plan: &'a PhysicalPlan,
245 write_op_node: PhysicalNodeId,
246 params: BTreeMap<String, LoraValue>,
247 ) -> ExecResult<Self> {
248 let input = match write_op_input(plan, write_op_node) {
249 Some(i) => i,
250 None => {
251 return Err(ExecutorError::RuntimeError(format!(
252 "StreamingWriteCursor::open called with non-write node {write_op_node:?}"
253 )));
254 }
255 };
256 let storage_ptr = StoragePtr::from_mut(storage);
257
258 let storage_ref: &'a S = unsafe { storage_ptr.as_ref() };
260 let upstream = build_streaming(plan, input, storage_ref, Arc::new(params.clone()))?;
261
262 Ok(Self {
263 upstream: ManuallyDrop::new(upstream),
264 storage_ptr,
265 plan,
266 write_op_node,
267 params,
268 _phantom: std::marker::PhantomData,
269 })
270 }
271}
272
273impl<'a, S: GraphStorageMut + GraphStorage + 'a> RowSource for StreamingWriteCursor<'a, S> {
274 fn next_row(&mut self) -> ExecResult<Option<Row>> {
275 let mut row = match self.upstream.next_row()? {
276 Some(r) => r,
277 None => return Ok(None),
278 };
279
280 let storage_mut: &mut S = unsafe { self.storage_ptr.as_mut() };
285 let mut exec = MutableExecutor::new(MutableExecutionContext {
286 storage: storage_mut,
287 params: self.params.clone(),
288 });
289 exec.defer_existence_checks(crate::executor::plan_defers_existence(self.plan));
292 let op = &self.plan.nodes[self.write_op_node];
293 exec.apply_write_op(op, &mut row)?;
294 exec.check_pending_existence()?;
295 let row = exec.hydrate_row(row);
296 Ok(Some(row))
297 }
298}
299
300impl<'a, S: GraphStorageMut + GraphStorage + 'a> Drop for StreamingWriteCursor<'a, S> {
301 fn drop(&mut self) {
302 unsafe {
306 ManuallyDrop::drop(&mut self.upstream);
307 }
308 }
309}
310
311struct DrainSilently<'a> {
314 inner: Option<Box<dyn RowSource + 'a>>,
315}
316
317impl RowSource for DrainSilently<'_> {
318 fn next_row(&mut self) -> ExecResult<Option<Row>> {
319 if let Some(mut inner) = self.inner.take() {
320 while inner.next_row()?.is_some() {}
321 }
322 Ok(None)
323 }
324}