uqa_execution/
project_set.rs1use uqa_sql::ResultRow;
10
11use crate::batch::DEFAULT_BATCH_SIZE;
12use crate::{Batch, ExecResult, OwnedPhysicalRow, PhysicalOperator, PhysicalRow, RowSchema};
13
14pub type ProjectRows = Box<dyn Iterator<Item = ExecResult<ResultRow>> + Send>;
21
22pub trait SetProjector: Send {
24 fn project(&mut self, row: &ResultRow) -> ExecResult<ProjectRows>;
25}
26
27impl<F> SetProjector for F
28where
29 F: FnMut(&ResultRow) -> ExecResult<ProjectRows> + Send,
30{
31 fn project(&mut self, row: &ResultRow) -> ExecResult<ProjectRows> {
32 self(row)
33 }
34}
35
36pub struct ProjectSet<'a> {
38 child: Box<dyn PhysicalOperator + 'a>,
39 projector: Box<dyn SetProjector + 'a>,
40 schema: RowSchema,
41 input: std::vec::IntoIter<ResultRow>,
42 projected: Option<ProjectRows>,
43 exhausted: bool,
44}
45
46impl<'a> ProjectSet<'a> {
47 pub fn new(
48 child: Box<dyn PhysicalOperator + 'a>,
49 output_schema: Vec<String>,
50 projector: Box<dyn SetProjector + 'a>,
51 ) -> Self {
52 Self {
53 child,
54 projector,
55 schema: RowSchema::new(output_schema),
56 input: Vec::new().into_iter(),
57 projected: None,
58 exhausted: false,
59 }
60 }
61
62 fn next_input(&mut self) -> ExecResult<Option<ResultRow>> {
63 loop {
64 if let Some(row) = self.input.next() {
65 return Ok(Some(row));
66 }
67 let Some(batch) = self.child.next()? else {
68 return Ok(None);
69 };
70 self.input = batch.into_result_rows().into_iter();
71 }
72 }
73}
74
75impl PhysicalOperator for ProjectSet<'_> {
76 fn row_schema(&self) -> &RowSchema {
77 &self.schema
78 }
79
80 fn open(&mut self) -> ExecResult<()> {
81 self.input = Vec::new().into_iter();
82 self.projected = None;
83 self.exhausted = false;
84 self.child.open()
85 }
86
87 fn next(&mut self) -> ExecResult<Option<Batch>> {
88 if self.exhausted && self.projected.is_none() {
89 return Ok(None);
90 }
91
92 let mut rows = Vec::with_capacity(DEFAULT_BATCH_SIZE);
93 while rows.len() < DEFAULT_BATCH_SIZE {
94 if let Some(projected) = self.projected.as_mut() {
95 match projected.next() {
96 Some(row) => {
97 rows.push(row?);
98 continue;
99 }
100 None => self.projected = None,
101 }
102 }
103
104 match self.next_input()? {
105 Some(row) => self.projected = Some(self.projector.project(&row)?),
106 None => {
107 self.exhausted = true;
108 break;
109 }
110 }
111 }
112 if rows.is_empty() {
113 return Ok(None);
114 }
115 Ok(Some(Batch::new(self.schema.clone(), rows)))
116 }
117
118 fn close(&mut self) -> ExecResult<()> {
119 self.input = Vec::new().into_iter();
120 self.projected = None;
121 self.exhausted = true;
122 self.child.close()
123 }
124}
125
126pub type PhysicalProjectRows = Box<dyn Iterator<Item = ExecResult<PhysicalRow>> + Send>;
128
129pub trait PhysicalSetProjector: Send {
131 fn project(&mut self, row: OwnedPhysicalRow) -> ExecResult<PhysicalProjectRows>;
132}
133
134impl<F> PhysicalSetProjector for F
135where
136 F: FnMut(OwnedPhysicalRow) -> ExecResult<PhysicalProjectRows> + Send,
137{
138 fn project(&mut self, row: OwnedPhysicalRow) -> ExecResult<PhysicalProjectRows> {
139 self(row)
140 }
141}
142
143pub struct PhysicalProjectSet<'a> {
145 child: Box<dyn PhysicalOperator + 'a>,
146 projector: Box<dyn PhysicalSetProjector + 'a>,
147 schema: RowSchema,
148 input: std::vec::IntoIter<OwnedPhysicalRow>,
149 projected: Option<PhysicalProjectRows>,
150 exhausted: bool,
151}
152
153impl<'a> PhysicalProjectSet<'a> {
154 pub fn new(
155 child: Box<dyn PhysicalOperator + 'a>,
156 schema: RowSchema,
157 projector: Box<dyn PhysicalSetProjector + 'a>,
158 ) -> Self {
159 Self {
160 child,
161 projector,
162 schema,
163 input: Vec::new().into_iter(),
164 projected: None,
165 exhausted: false,
166 }
167 }
168
169 fn next_input(&mut self) -> ExecResult<Option<OwnedPhysicalRow>> {
170 loop {
171 if let Some(row) = self.input.next() {
172 return Ok(Some(row));
173 }
174 let Some(batch) = self.child.next()? else {
175 return Ok(None);
176 };
177 self.input = batch.into_owned_rows().into_iter();
178 }
179 }
180}
181
182impl PhysicalOperator for PhysicalProjectSet<'_> {
183 fn row_schema(&self) -> &RowSchema {
184 &self.schema
185 }
186
187 fn open(&mut self) -> ExecResult<()> {
188 self.input = Vec::new().into_iter();
189 self.projected = None;
190 self.exhausted = false;
191 self.child.open()
192 }
193
194 fn next(&mut self) -> ExecResult<Option<Batch>> {
195 if self.exhausted && self.projected.is_none() {
196 return Ok(None);
197 }
198
199 let mut rows = Vec::with_capacity(DEFAULT_BATCH_SIZE);
200 while rows.len() < DEFAULT_BATCH_SIZE {
201 if let Some(projected) = self.projected.as_mut() {
202 match projected.next() {
203 Some(row) => {
204 rows.push(row?);
205 continue;
206 }
207 None => self.projected = None,
208 }
209 }
210
211 match self.next_input()? {
212 Some(row) => self.projected = Some(self.projector.project(row)?),
213 None => {
214 self.exhausted = true;
215 break;
216 }
217 }
218 }
219 if rows.is_empty() {
220 return Ok(None);
221 }
222 Ok(Some(Batch::from_physical_rows(self.schema.clone(), rows)))
223 }
224
225 fn close(&mut self) -> ExecResult<()> {
226 self.input = Vec::new().into_iter();
227 self.projected = None;
228 self.exhausted = true;
229 self.child.close()
230 }
231}
232
233#[cfg(test)]
234mod tests {
235 use super::*;
236 use crate::physical::run_to_rows;
237 use crate::scan::TableScan;
238 use crate::{RowLockOrigin, RowProjectionValue};
239 use uqa_core::Value;
240
241 #[test]
242 fn expands_each_input_row_in_child_order() {
243 let input: ResultRow = [("n".into(), Value::Int(2))].into_iter().collect();
244 let child = TableScan::from_rows(vec!["n".into()], vec![input]);
245 let projector = |row: &ResultRow| -> ExecResult<ProjectRows> {
246 let Value::Int(end) = row.get("n").cloned().unwrap_or(Value::Null) else {
247 return Ok(Box::new(std::iter::empty()));
248 };
249 Ok(Box::new((1..=end).map(|value| {
250 Ok([("value".into(), Value::Int(value))].into_iter().collect())
251 })))
252 };
253 let mut project =
254 ProjectSet::new(Box::new(child), vec!["value".into()], Box::new(projector));
255 let (_, rows) = run_to_rows(&mut project).unwrap();
256 assert_eq!(rows.len(), 2);
257 assert_eq!(rows[0].get("value"), Some(&Value::Int(1)));
258 assert_eq!(rows[1].get("value"), Some(&Value::Int(2)));
259 }
260
261 #[test]
262 fn one_input_with_many_outputs_never_builds_an_unbounded_pending_queue() {
263 let input: ResultRow = [("n".into(), Value::Int(10_000))].into_iter().collect();
264 let child = TableScan::from_rows(vec!["n".into()], vec![input]);
265 let projector = |row: &ResultRow| -> ExecResult<ProjectRows> {
266 let Value::Int(end) = row.get("n").cloned().unwrap_or(Value::Null) else {
267 return Ok(Box::new(std::iter::empty()));
268 };
269 Ok(Box::new((0..end).map(|value| {
270 Ok([("value".into(), Value::Int(value))].into_iter().collect())
271 })))
272 };
273 let mut project =
274 ProjectSet::new(Box::new(child), vec!["value".into()], Box::new(projector));
275 project.open().unwrap();
276 let first = project.next().unwrap().unwrap();
277 assert_eq!(first.len(), DEFAULT_BATCH_SIZE);
278 let second = project.next().unwrap().unwrap();
279 assert_eq!(second.len(), DEFAULT_BATCH_SIZE);
280 project.close().unwrap();
281 }
282
283 #[test]
284 fn late_projector_error_is_not_converted_to_end_of_stream() {
285 let child = TableScan::from_rows(vec!["n".into()], vec![ResultRow::new()]);
286 let projector = |_row: &ResultRow| -> ExecResult<ProjectRows> {
287 Ok(Box::new(
288 vec![
289 Ok([("value".into(), Value::Int(1))].into_iter().collect()),
290 Err(crate::ExecError::Other("injected projector failure".into())),
291 ]
292 .into_iter(),
293 ))
294 };
295 let mut project =
296 ProjectSet::new(Box::new(child), vec!["value".into()], Box::new(projector));
297 project.open().unwrap();
298 let error = project.next().unwrap_err();
299 assert!(error.to_string().contains("injected projector failure"));
300 }
301
302 #[test]
303 fn physical_projection_preserves_row_lineage_without_named_materialization() {
304 let input_schema = RowSchema::new(vec!["unused".into(), "value".into()]);
305 let input_row = PhysicalRow::from_values(vec![
306 Value::Str("unused payload".repeat(64)),
307 Value::Str("shared value".repeat(64)),
308 ])
309 .with_lock_origin(RowLockOrigin::new("source", "public.source", 7));
310 let child = TableScan::from_physical_rows(input_schema, vec![input_row]);
311 let projector = |row: OwnedPhysicalRow| -> ExecResult<PhysicalProjectRows> {
312 let slot = row.schema.physical_slot(1).unwrap();
313 Ok(Box::new(std::iter::once(Ok(row.row.project_with_values(
314 [RowProjectionValue::InputSlot(slot)],
315 )))))
316 };
317 let output_schema = RowSchema::new(vec!["value".into()]);
318 let mut project =
319 PhysicalProjectSet::new(Box::new(child), output_schema.clone(), Box::new(projector));
320
321 project.open().unwrap();
322 let output = project.next().unwrap().unwrap();
323 assert_eq!(output.rows.len(), 1);
324 assert_eq!(
325 output_schema.view(&output.rows[0]).value_at(0),
326 Some(&Value::Str("shared value".repeat(64)))
327 );
328 assert_eq!(output.rows[0].lock_origins()[0].doc_id, 7);
329 assert!(project.next().unwrap().is_none());
330 project.close().unwrap();
331 }
332}