Skip to main content

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}