Skip to main content

uqa_execution/
project_set.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Streaming one-to-many projection (`ProjectSet`).
8
9use uqa_sql::ResultRow;
10
11use crate::batch::DEFAULT_BATCH_SIZE;
12use crate::{Batch, ExecResult, OwnedPhysicalRow, PhysicalOperator, PhysicalRow, RowSchema};
13
14/// Owned row stream produced for one input row.
15///
16/// The iterator yields errors at the exact row where production failed. This
17/// matters for generators backed by files or user code: eagerly collecting a
18/// `Vec` both defeated the Volcano memory boundary and made late failures
19/// impossible to represent without discarding already-produced rows.
20pub type ProjectRows = Box<dyn Iterator<Item = ExecResult<ResultRow>> + Send>;
21
22/// Engine seam for set-returning projection expressions.
23pub 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
36/// Pulls input rows and expands each through a set-returning projection.
37pub 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
126/// Owned physical row stream produced for one input row without materializing named maps.
127pub type PhysicalProjectRows = Box<dyn Iterator<Item = ExecResult<PhysicalRow>> + Send>;
128
129/// Engine seam for set-returning projections that keep rows in their physical layout.
130pub 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
143/// Pulls physical input rows and expands each without crossing the named-row materialization boundary.
144pub 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}