Skip to main content

uqa_execution/spill/
indexed.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7use std::fs::File;
8use std::io::{Read, Seek, SeekFrom, Write};
9
10use tempfile::NamedTempFile;
11
12use crate::batch::{PhysicalRow, RowSchema};
13use crate::physical::ExecResult;
14
15use super::format::{
16    decode_physical_row_record, encode_physical_row_record, encoded_physical_row_record_size,
17    spill_error, RECORD_PREFIX_BYTES,
18};
19
20#[cfg(test)]
21mod tests;
22
23/// Disk-only physical-row store with constant-memory positional lookup.
24///
25/// Each row retains the exact physical layout described by `schema` and is
26/// encoded with the same positional tagged representation as [`super::SpillBuffer`].
27/// Record offsets live in a second temporary file, so even a partition with
28/// billions of rows does not create an in-memory offset table. A single decoded
29/// physical row is the only input-sized allocation retained by [`Self::get`].
30/// Both files are unlinked by `NamedTempFile` on drop.
31pub struct IndexedSpill {
32    schema: RowSchema,
33    data: NamedTempFile,
34    offsets: NamedTempFile,
35    rows: u64,
36    encoded_bytes: u64,
37}
38
39impl IndexedSpill {
40    pub fn new(input_schema: RowSchema) -> ExecResult<Self> {
41        Ok(Self {
42            schema: input_schema,
43            data: NamedTempFile::new().map_err(|error| {
44                spill_error(format!("failed to create indexed spill data: {error}"))
45            })?,
46            offsets: NamedTempFile::new().map_err(|error| {
47                spill_error(format!("failed to create indexed spill offsets: {error}"))
48            })?,
49            rows: 0,
50            encoded_bytes: 0,
51        })
52    }
53
54    pub fn len(&self) -> u64 {
55        self.rows
56    }
57
58    pub fn is_empty(&self) -> bool {
59        self.rows == 0
60    }
61
62    pub fn encoded_bytes(&self) -> u64 {
63        self.encoded_bytes
64    }
65
66    pub fn row_schema(&self) -> &RowSchema {
67        &self.schema
68    }
69
70    pub(crate) fn encoded_row_size(schema: &RowSchema, row: &PhysicalRow) -> ExecResult<usize> {
71        encoded_physical_row_record_size(row, schema.physical_width())?
72            .checked_add(RECORD_PREFIX_BYTES)
73            .ok_or_else(|| spill_error("indexed spill row size overflow"))
74    }
75
76    /// Append one indivisible row. Failed writes roll both files back to their
77    /// original lengths, so callers never observe a partial index entry.
78    pub fn push(&mut self, row: &PhysicalRow) -> ExecResult<()> {
79        let payload = encode_physical_row_record(row, self.schema.physical_width())?;
80        let length = u64::try_from(payload.len())
81            .map_err(|_| spill_error("indexed spill row is too large"))?;
82        // Validate every piece of metadata before touching either file.  A
83        // counter overflow after the append would otherwise return an error
84        // while leaving a physically visible row whose offset/count was not
85        // published consistently.
86        let next_rows = self
87            .rows
88            .checked_add(1)
89            .ok_or_else(|| spill_error("indexed spill row count overflow"))?;
90        let record_bytes = length
91            .checked_add(8)
92            .ok_or_else(|| spill_error("indexed spill row length overflow"))?;
93        let next_encoded_bytes = self
94            .encoded_bytes
95            .checked_add(record_bytes)
96            .ok_or_else(|| spill_error("indexed spill byte count overflow"))?;
97        let data_length = self
98            .data
99            .as_file_mut()
100            .seek(SeekFrom::End(0))
101            .map_err(|error| spill_error(format!("failed to seek indexed spill data: {error}")))?;
102        let offsets_length =
103            self.offsets
104                .as_file_mut()
105                .seek(SeekFrom::End(0))
106                .map_err(|error| {
107                    spill_error(format!("failed to seek indexed spill offsets: {error}"))
108                })?;
109
110        let write_result = (|| -> std::io::Result<()> {
111            self.data.as_file_mut().write_all(&length.to_le_bytes())?;
112            self.data.as_file_mut().write_all(&payload)?;
113            self.data.as_file_mut().flush()?;
114            self.offsets
115                .as_file_mut()
116                .write_all(&data_length.to_le_bytes())?;
117            self.offsets.as_file_mut().flush()
118        })();
119        if let Err(error) = write_result {
120            let data_rollback = self.data.as_file_mut().set_len(data_length);
121            let offsets_rollback = self.offsets.as_file_mut().set_len(offsets_length);
122            let rollback_error = match (data_rollback, offsets_rollback) {
123                (Ok(()), Ok(())) => None,
124                (Err(data), Ok(())) => Some(format!("data rollback failed: {data}")),
125                (Ok(()), Err(offsets)) => Some(format!("offset rollback failed: {offsets}")),
126                (Err(data), Err(offsets)) => Some(format!(
127                    "data rollback failed: {data}; offset rollback failed: {offsets}"
128                )),
129            };
130            if let Some(rollback) = rollback_error {
131                return Err(spill_error(format!(
132                    "failed to append indexed spill row: {error}; {rollback}"
133                )));
134            }
135            return Err(spill_error(format!(
136                "failed to append indexed spill row: {error}"
137            )));
138        }
139
140        self.rows = next_rows;
141        self.encoded_bytes = next_encoded_bytes;
142        Ok(())
143    }
144
145    /// Decode the row at `index` without loading any other row or index entry.
146    pub fn get(&mut self, index: u64) -> ExecResult<PhysicalRow> {
147        if index >= self.rows {
148            return Err(spill_error(format!(
149                "indexed spill row {index} is outside 0..{}",
150                self.rows
151            )));
152        }
153        let expected_offsets_length = self
154            .rows
155            .checked_mul(8)
156            .ok_or_else(|| spill_error("indexed spill offsets length overflow"))?;
157        let actual_offsets_length = self
158            .offsets
159            .as_file()
160            .metadata()
161            .map_err(|error| {
162                spill_error(format!("failed to inspect indexed spill offsets: {error}"))
163            })?
164            .len();
165        if actual_offsets_length != expected_offsets_length {
166            return Err(spill_error(format!(
167                "indexed spill offsets length {actual_offsets_length} does not match expected {expected_offsets_length}"
168            )));
169        }
170        let data_length = self
171            .data
172            .as_file()
173            .metadata()
174            .map_err(|error| spill_error(format!("failed to inspect indexed spill data: {error}")))?
175            .len();
176        let offset_position = index
177            .checked_mul(8)
178            .ok_or_else(|| spill_error("indexed spill offset overflow"))?;
179        let offset = read_indexed_offset(self.offsets.as_file_mut(), offset_position)?;
180        let record_end = if index
181            .checked_add(1)
182            .ok_or_else(|| spill_error("indexed spill row index overflow"))?
183            < self.rows
184        {
185            read_indexed_offset(
186                self.offsets.as_file_mut(),
187                offset_position
188                    .checked_add(8)
189                    .ok_or_else(|| spill_error("indexed spill next offset overflow"))?,
190            )?
191        } else {
192            data_length
193        };
194        let payload_start = offset
195            .checked_add(8)
196            .ok_or_else(|| spill_error("indexed spill payload offset overflow"))?;
197        if payload_start > record_end || record_end > data_length {
198            return Err(spill_error(format!(
199                "indexed spill record bounds {offset}..{record_end} are outside data length {data_length}"
200            )));
201        }
202        self.data
203            .as_file_mut()
204            .seek(SeekFrom::Start(offset))
205            .map_err(|error| spill_error(format!("failed to seek indexed spill row: {error}")))?;
206        let mut length = [0_u8; 8];
207        self.data
208            .as_file_mut()
209            .read_exact(&mut length)
210            .map_err(|error| {
211                spill_error(format!("failed to read indexed spill length: {error}"))
212            })?;
213        let declared_length = u64::from_le_bytes(length);
214        let available_length = record_end - payload_start;
215        if declared_length != available_length {
216            return Err(spill_error(format!(
217                "indexed spill row length {declared_length} does not match record payload {available_length}"
218            )));
219        }
220        let length = usize::try_from(declared_length)
221            .map_err(|_| spill_error("indexed spill row length is outside address space"))?;
222        let mut payload = Vec::new();
223        payload.try_reserve_exact(length).map_err(|error| {
224            spill_error(format!(
225                "unable to allocate indexed spill row payload of {length} bytes: {error}"
226            ))
227        })?;
228        payload.resize(length, 0);
229        self.data
230            .as_file_mut()
231            .read_exact(&mut payload)
232            .map_err(|error| spill_error(format!("failed to read indexed spill row: {error}")))?;
233        decode_physical_row_record(&payload, self.schema.physical_width())
234    }
235}
236
237fn read_indexed_offset(file: &mut File, position: u64) -> ExecResult<u64> {
238    file.seek(SeekFrom::Start(position))
239        .map_err(|error| spill_error(format!("failed to seek indexed spill offset: {error}")))?;
240    let mut encoded = [0_u8; 8];
241    file.read_exact(&mut encoded)
242        .map_err(|error| spill_error(format!("failed to read indexed spill offset: {error}")))?;
243    Ok(u64::from_le_bytes(encoded))
244}