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 )
33 })
34}
35
36/// Classify a compiled query for streaming. Treats any UNION branch
37/// the same as the head: a single mutating op anywhere across the
38/// compiled query promotes the whole query to `Mutating`.
39pub fn classify_stream(compiled: &CompiledQuery) -> StreamShape {
40 if plan_is_mutating(&compiled.physical)
41 || compiled
42 .unions
43 .iter()
44 .any(|b| plan_is_mutating(&b.physical))
45 {
46 StreamShape::Mutating
47 } else {
48 StreamShape::ReadOnly
49 }
50}