uqa_execution/spill/
indexed.rs1use 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
23pub 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 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 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 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}