Skip to main content

uqa_execution/
scroll_materialize.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Incremental directional materialization for scrollable executor boundaries.
8
9use std::collections::VecDeque;
10
11use crate::{
12    BackwardScanSupport, Batch, ExecError, ExecResult, IndexedSpill, PhysicalOperator,
13    PhysicalOrder, PhysicalRow, PhysicalScanDirection, RowSchema,
14};
15
16#[derive(Clone, Copy)]
17enum MaterializePosition {
18    BeforeFirst,
19    OnRow(u64),
20    AfterLast,
21}
22
23/// Cache an operator's output as it is first pulled forward, then expose the cached rows in either direction. This is the physical equivalent of `PostgreSQL`'s `Material` node: expressions below the boundary run once, while expressions in parents run again when a row is revisited.
24pub struct ScrollMaterialize<'a> {
25    child: Box<dyn PhysicalOperator + 'a>,
26    schema: RowSchema,
27    ordering: Vec<PhysicalOrder>,
28    rows: Option<IndexedSpill>,
29    pending: VecDeque<PhysicalRow>,
30    position: MaterializePosition,
31    eof: bool,
32}
33
34impl<'a> ScrollMaterialize<'a> {
35    pub fn new(child: Box<dyn PhysicalOperator + 'a>) -> Self {
36        let schema = child.row_schema().clone();
37        let ordering = child.output_ordering().to_vec();
38        Self {
39            child,
40            schema,
41            ordering,
42            rows: None,
43            pending: VecDeque::new(),
44            position: MaterializePosition::BeforeFirst,
45            eof: false,
46        }
47    }
48
49    fn rows_mut(&mut self) -> ExecResult<&mut IndexedSpill> {
50        self.rows
51            .as_mut()
52            .ok_or_else(|| ExecError::Other("scroll materialization is not open".into()))
53    }
54
55    fn pull_child_row(&mut self) -> ExecResult<Option<PhysicalRow>> {
56        loop {
57            if let Some(row) = self.pending.pop_front() {
58                return Ok(Some(row));
59            }
60            let Some(batch) = self.child.next()? else {
61                return Ok(None);
62            };
63            if batch.schema != self.schema {
64                return Err(ExecError::Other(format!(
65                    "scroll materialization input schema mismatch: expected {:?}, got {:?}",
66                    self.schema, batch.schema
67                )));
68            }
69            self.pending.extend(batch.rows);
70        }
71    }
72
73    fn row_at(&mut self, index: u64) -> ExecResult<Batch> {
74        let row = self.rows_mut()?.get(index)?;
75        Ok(Batch::from_physical_rows(self.schema.clone(), vec![row]))
76    }
77
78    fn next_forward(&mut self) -> ExecResult<Option<Batch>> {
79        let target = match self.position {
80            MaterializePosition::BeforeFirst => 0,
81            MaterializePosition::OnRow(position) => position
82                .checked_add(1)
83                .ok_or_else(|| ExecError::Other("scroll position overflow".into()))?,
84            MaterializePosition::AfterLast => return Ok(None),
85        };
86        if target < self.rows_mut()?.len() {
87            self.position = MaterializePosition::OnRow(target);
88            return self.row_at(target).map(Some);
89        }
90        if self.eof {
91            self.position = MaterializePosition::AfterLast;
92            return Ok(None);
93        }
94        let Some(row) = self.pull_child_row()? else {
95            self.eof = true;
96            self.position = MaterializePosition::AfterLast;
97            return Ok(None);
98        };
99        self.rows_mut()?.push(&row)?;
100        self.position = MaterializePosition::OnRow(target);
101        Ok(Some(Batch::from_physical_rows(
102            self.schema.clone(),
103            vec![row],
104        )))
105    }
106
107    fn next_backward(&mut self) -> ExecResult<Option<Batch>> {
108        let target = match self.position {
109            MaterializePosition::BeforeFirst => return Ok(None),
110            MaterializePosition::OnRow(0) => {
111                self.position = MaterializePosition::BeforeFirst;
112                return Ok(None);
113            }
114            MaterializePosition::OnRow(position) => position - 1,
115            MaterializePosition::AfterLast => {
116                let row_count = self.rows_mut()?.len();
117                let Some(position) = row_count.checked_sub(1) else {
118                    self.position = MaterializePosition::BeforeFirst;
119                    return Ok(None);
120                };
121                position
122            }
123        };
124        self.position = MaterializePosition::OnRow(target);
125        self.row_at(target).map(Some)
126    }
127}
128
129impl PhysicalOperator for ScrollMaterialize<'_> {
130    fn row_schema(&self) -> &RowSchema {
131        &self.schema
132    }
133
134    fn estimated_cardinality(&self) -> Option<u64> {
135        self.child.estimated_cardinality()
136    }
137
138    fn output_ordering(&self) -> &[PhysicalOrder] {
139        &self.ordering
140    }
141
142    fn backward_scan_support(&self) -> BackwardScanSupport {
143        BackwardScanSupport::Native
144    }
145
146    fn open(&mut self) -> ExecResult<()> {
147        self.pending.clear();
148        self.position = MaterializePosition::BeforeFirst;
149        self.eof = false;
150        self.rows = Some(IndexedSpill::new(self.schema.clone())?);
151        self.child.open()
152    }
153
154    fn next(&mut self) -> ExecResult<Option<Batch>> {
155        self.next_forward()
156    }
157
158    fn next_direction(&mut self, direction: PhysicalScanDirection) -> ExecResult<Option<Batch>> {
159        match direction {
160            PhysicalScanDirection::Forward => self.next_forward(),
161            PhysicalScanDirection::Backward => self.next_backward(),
162        }
163    }
164
165    fn rewind(&mut self) -> ExecResult<()> {
166        self.position = MaterializePosition::BeforeFirst;
167        Ok(())
168    }
169
170    fn close(&mut self) -> ExecResult<()> {
171        self.pending.clear();
172        self.rows = None;
173        self.position = MaterializePosition::BeforeFirst;
174        self.eof = true;
175        self.child.close()
176    }
177}
178
179/// Materialize an operator only when it identifies its output as a safe semantic boundary.
180pub fn prepare_backward_scan<'a>(
181    operator: Box<dyn PhysicalOperator + 'a>,
182) -> Box<dyn PhysicalOperator + 'a> {
183    if operator.backward_scan_support() == BackwardScanSupport::Materialize {
184        Box::new(ScrollMaterialize::new(operator))
185    } else {
186        operator
187    }
188}