Skip to main content

uqa_execution/
spill_scan.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Volcano scan over an owned [`crate::spill::SpillBuffer`].
8//!
9//! The scan transfers ownership of the spill file to its batch iterator at
10//! `open`, so disk batches are decoded one at a time and the file is removed
11//! even when execution stops early.
12
13use crate::batch::{Batch, RowSchema};
14use crate::physical::{BackwardScanSupport, ExecResult, PhysicalOperator};
15use crate::spill::{SharedSpill, SharedSpillReader, SpillBuffer, SpillDrain};
16
17/// One-shot physical scan over a disk-backed spill buffer.
18pub struct SpillScan {
19    schema: RowSchema,
20    buffer: Option<SpillBuffer>,
21    reader: Option<SpillDrain>,
22}
23
24/// Repeatable scan over an immutable shared spill. Cloning the source only
25/// clones an `Arc`; `open` creates an independent file reader.
26pub struct SharedSpillScan {
27    source: SharedSpill,
28    schema: RowSchema,
29    reader: Option<SharedSpillReader>,
30}
31
32impl SharedSpillScan {
33    pub fn new(source: SharedSpill) -> Self {
34        let schema = source.row_schema().clone();
35        Self {
36            source,
37            schema,
38            reader: None,
39        }
40    }
41}
42
43impl PhysicalOperator for SharedSpillScan {
44    fn row_schema(&self) -> &RowSchema {
45        &self.schema
46    }
47
48    fn estimated_cardinality(&self) -> Option<u64> {
49        u64::try_from(self.source.rows()).ok()
50    }
51
52    fn backward_scan_support(&self) -> BackwardScanSupport {
53        BackwardScanSupport::Materialize
54    }
55
56    fn open(&mut self) -> ExecResult<()> {
57        self.reader = Some(self.source.reader()?);
58        Ok(())
59    }
60
61    fn next(&mut self) -> ExecResult<Option<Batch>> {
62        self.reader
63            .as_mut()
64            .map_or(Ok(None), |reader| reader.next().transpose())
65    }
66
67    fn close(&mut self) -> ExecResult<()> {
68        self.reader = None;
69        Ok(())
70    }
71}
72
73impl SpillScan {
74    pub fn new(schema: impl Into<RowSchema>, buffer: SpillBuffer) -> Self {
75        let schema = schema.into();
76        Self {
77            schema,
78            buffer: Some(buffer),
79            reader: None,
80        }
81    }
82}
83
84impl PhysicalOperator for SpillScan {
85    fn row_schema(&self) -> &RowSchema {
86        &self.schema
87    }
88
89    fn open(&mut self) -> ExecResult<()> {
90        let mut buffer = self.buffer.take().ok_or_else(|| {
91            crate::physical::ExecError::Other("spill scan cannot be reopened".into())
92        })?;
93        self.reader = Some(buffer.drain()?);
94        Ok(())
95    }
96
97    fn next(&mut self) -> ExecResult<Option<Batch>> {
98        let Some(reader) = self.reader.as_mut() else {
99            return Ok(None);
100        };
101        let Some(batch) = reader.next().transpose()? else {
102            return Ok(None);
103        };
104        if batch.schema != self.schema {
105            return Err(crate::physical::ExecError::Other(format!(
106                "spill scan schema mismatch: expected {:?}, got {:?}",
107                self.schema.columns(),
108                batch.schema.columns()
109            )));
110        }
111        Ok(Some(batch))
112    }
113
114    fn close(&mut self) -> ExecResult<()> {
115        self.reader = None;
116        Ok(())
117    }
118}
119
120#[cfg(test)]
121mod tests {
122    use std::collections::BTreeMap;
123
124    use uqa_core::Value;
125
126    use super::*;
127    use crate::physical::run_to_rows;
128
129    #[test]
130    fn scan_streams_a_forced_spill_in_input_order() {
131        let schema = RowSchema::new(vec!["x".into()]);
132        let mut spill = SpillBuffer::new(1);
133        for value in 0..300_i64 {
134            spill
135                .push(Batch::new(
136                    schema.clone(),
137                    vec![BTreeMap::from([("x".into(), Value::Int(value))])],
138                ))
139                .unwrap();
140        }
141        assert!(spill.has_spilled());
142
143        let mut scan = SpillScan::new(schema.columns().to_vec(), spill);
144        let (_, rows) = run_to_rows(&mut scan).unwrap();
145        assert_eq!(rows.len(), 300);
146        for (expected, row) in rows.iter().enumerate() {
147            assert_eq!(row.get("x"), Some(&Value::Int(expected as i64)));
148        }
149    }
150
151    #[test]
152    fn shared_spill_supports_independent_repeatable_scans() {
153        let schema = RowSchema::new(vec!["x".into()]);
154        let mut spill = SpillBuffer::new(1);
155        for value in 0..2_048_i64 {
156            spill
157                .push(Batch::new(
158                    schema.clone(),
159                    vec![BTreeMap::from([("x".into(), Value::Int(value))])],
160                ))
161                .unwrap();
162        }
163        let shared = spill.into_shared(schema.columns().to_vec()).unwrap();
164        let mut first = SharedSpillScan::new(shared.clone());
165        let mut second = SharedSpillScan::new(shared);
166        let (_, first_rows) = run_to_rows(&mut first).unwrap();
167        let (_, second_rows) = run_to_rows(&mut second).unwrap();
168        assert_eq!(first_rows, second_rows);
169        assert_eq!(first_rows.len(), 2_048);
170    }
171}