Skip to main content

lora_executor/pull/
mutable.rs

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
16/// Pull-based read-write executor. Wraps the existing
17/// [`MutableExecutor`] under the same row-cursor API. Mutations are
18/// applied during `open_compiled`; the returned cursor yields the
19/// resulting rows lazily.
20pub 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    /// Open a cursor for a compiled write query.
31    ///
32    /// Fast path: when a branch root is one of `Create` / `Set` /
33    /// `Delete` / `Remove` / `Merge` and its input subtree is fully
34    /// streamable, returns a [`StreamingWriteCursor`] that pulls input
35    /// row-by-row and applies the per-row write through
36    /// [`MutableExecutor::apply_write_op`]. `UNION ALL` plans stream
37    /// one branch at a time. Plain `UNION` drains branches first so
38    /// rows can be deduplicated by name.
39    ///
40    /// Fallback: a branch that is not streamable materializes through
41    /// [`MutableExecutor::execute_rows`] and wraps the result in a
42    /// [`BufferedRowSource`].
43    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
100/// Mutable UNION cursor. `UNION ALL` streams one branch at a time
101/// against the same staged graph. Plain `UNION` streams branch-by-branch
102/// while retaining only a seen-key set for deduplication.
103pub 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        // SAFETY: MutableUnionSource keeps at most one branch cursor
148        // alive at a time. `current` is dropped before advancing to
149        // the next branch, so each mutable reborrow is temporally
150        // disjoint.
151        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
197/// Streaming write cursor for plans whose root is one of
198/// `Create` / `Set` / `Delete` / `Remove` / `Merge` and whose input
199/// subtree is fully streamable.
200///
201/// # Layout invariant
202///
203/// The cursor owns a raw alias of the original `&'a mut S`.
204/// Its `upstream` was constructed using a `&'a S` reborrow derived
205/// from `storage_ptr` via unsafe lifetime extension. This is sound
206/// because the existing read-side `RowSource` impls (see
207/// `NodeScanSource::cur_ids`, `ExpandSource::cur_edges`, etc.)
208/// materialize their iteration state into owned `Vec`s at
209/// construction or first call, so no live `&S` borrow into storage
210/// persists across `next_row` calls. Read-only access happens
211/// transiently inside each `upstream.next_row` call; mutable access
212/// happens between calls inside [`MutableExecutor::apply_write_op`].
213/// The borrows never overlap in time.
214///
215/// # Drop order
216///
217/// `upstream` must drop before any caller may regain `&mut S` access
218/// to the underlying storage. The explicit `Drop` impl enforces
219/// that order — `ManuallyDrop` lets us force the sequence.
220pub struct StreamingWriteCursor<'a, S: GraphStorageMut + GraphStorage + 'a> {
221    /// SAFETY: borrows from `*storage_ptr`. Must drop first.
222    upstream: ManuallyDrop<Box<dyn RowSource + 'a>>,
223    /// Raw alias of the `&'a mut S` handed in at construction. Used
224    /// as `&S` by `upstream` and as `&mut S` inside this cursor's `next_row`.
225    storage_ptr: StoragePtr<S>,
226    /// Physical plan — kept alive for the per-row op borrow.
227    plan: &'a PhysicalPlan,
228    /// Index into `plan.nodes` of the write operator.
229    /// We re-fetch the op per call so this struct doesn't need to
230    /// be parameterized by the specific op type.
231    write_op_node: PhysicalNodeId,
232    /// Parameters; cloned per row into a fresh `MutableExecutor`.
233    /// In typical bulk-write workloads this is empty or tiny.
234    params: BTreeMap<String, LoraValue>,
235    _phantom: std::marker::PhantomData<&'a mut S>,
236}
237
238impl<'a, S: GraphStorageMut + GraphStorage + 'a> StreamingWriteCursor<'a, S> {
239    /// Build a cursor. Caller must already have verified that
240    /// `plan.nodes[write_op_node]` is a streamable write op via
241    /// [`write_op_input`] and [`subtree_is_fully_streaming`].
242    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        // SAFETY: see struct-level comment.
259        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        // SAFETY: upstream's `next_row` has returned, so its
281        // dormant `&S` borrow is not in active use right now. We
282        // reborrow `&mut S` for the per-row write and drop the
283        // borrow before the next pull.
284        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        // The write op is the only one in the plan, so a row is a whole
290        // statement's worth of writes for existence checks.
291        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        // SAFETY: drop `upstream` first to release its borrow into
303        // `*storage_ptr`. Subsequent fields drop via the normal
304        // field-drop sequence and don't touch storage.
305        unsafe {
306            ManuallyDrop::drop(&mut self.upstream);
307        }
308    }
309}
310
311/// Runs a write statement without `RETURN` to completion on the first pull
312/// and yields no rows.
313struct 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}