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