lora_executor/pull/shape.rs
1use lora_compiler::physical::{PhysicalOp, PhysicalPlan};
2use lora_compiler::CompiledQuery;
3
4/// Classification of a compiled query, used by the database layer to
5/// decide whether `db.stream` needs a hidden staged transaction.
6#[derive(Debug, Clone, Copy, PartialEq, Eq)]
7pub enum StreamShape {
8 /// No mutating operator anywhere in the plan or any of its
9 /// UNION branches. Safe to stream against the live store.
10 ReadOnly,
11 /// Has at least one mutating operator (Create / Merge / Delete /
12 /// Set / Remove). The host should run this against a staged
13 /// graph and only publish on cursor exhaustion.
14 Mutating,
15}
16
17impl StreamShape {
18 pub fn is_mutating(self) -> bool {
19 matches!(self, StreamShape::Mutating)
20 }
21}
22
23fn plan_is_mutating(plan: &PhysicalPlan) -> bool {
24 plan.nodes.iter().any(|op| {
25 matches!(
26 op,
27 PhysicalOp::Create(_)
28 | PhysicalOp::Merge(_)
29 | PhysicalOp::Delete(_)
30 | PhysicalOp::Set(_)
31 | PhysicalOp::Remove(_)
32 | PhysicalOp::Foreach(_)
33 )
34 })
35}
36
37/// Classify a compiled query for streaming. Treats any UNION branch
38/// the same as the head: a single mutating op anywhere across the
39/// compiled query promotes the whole query to `Mutating`.
40pub fn classify_stream(compiled: &CompiledQuery) -> StreamShape {
41 if plan_is_mutating(&compiled.physical)
42 || compiled
43 .unions
44 .iter()
45 .any(|b| plan_is_mutating(&b.physical))
46 {
47 StreamShape::Mutating
48 } else {
49 StreamShape::ReadOnly
50 }
51}