uqa_execution/
scroll_materialize.rs1use 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
23pub 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
179pub 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}