1use std::cmp::Ordering;
18use std::path::PathBuf;
19
20use uqa_core::Value;
21
22use crate::batch::{Batch, PhysicalRow, RowSchema};
23use crate::physical::{
24 order_expression_position, ExecError, ExecResult, PhysicalOperator, PhysicalOrder,
25};
26use crate::relational::{compare_sort_key_values_by, SharedExpressionEvaluator, SortKey};
27use crate::spill::{EncodedBatchSizer, SpillBuffer, SpillDrain};
28
29pub const EXTERNAL_SORT_MERGE_FAN_IN: usize = 16;
31
32const RUN_SOURCE_PREFIX: &str = "\0uqa.external_sort.source.";
33const RUN_KEY_PREFIX: &str = "\0uqa.external_sort.key.";
34const RUN_SEQUENCE: &str = "\0uqa.external_sort.sequence";
35
36fn run_schema(source_width: usize, key_count: usize) -> RowSchema {
37 let mut columns = (0..source_width)
38 .map(|index| format!("{RUN_SOURCE_PREFIX}{index}"))
39 .collect::<Vec<_>>();
40 columns.extend((0..key_count).map(|index| format!("{RUN_KEY_PREFIX}{index}")));
41 columns.push(RUN_SEQUENCE.to_string());
42 RowSchema::new(columns)
43}
44
45pub struct ExternalSort<'a> {
47 child: Box<dyn PhysicalOperator + 'a>,
48 keys: Vec<SortKey>,
49 evaluator: SharedExpressionEvaluator<'a>,
50 keep: Option<usize>,
51 work_mem_bytes: usize,
52 spill_directory: Option<PathBuf>,
53 schema: RowSchema,
54 input_slots: Vec<usize>,
55 run_schema: RowSchema,
56 ordering: Vec<PhysicalOrder>,
57 output: Option<SpillDrain>,
58 initial_run_count: usize,
59 merge_pass_count: usize,
60}
61
62impl<'a> ExternalSort<'a> {
63 pub fn new(
64 child: Box<dyn PhysicalOperator + 'a>,
65 keys: Vec<SortKey>,
66 evaluator: SharedExpressionEvaluator<'a>,
67 keep: Option<usize>,
68 work_mem_bytes: usize,
69 ) -> Self {
70 let (schema, input_slots) = child.row_schema().canonical_projection();
71 let run_schema = run_schema(input_slots.len(), keys.len());
72 let ordering = keys
73 .iter()
74 .map(|key| {
75 order_expression_position(&schema, &key.expr).map(|position| PhysicalOrder {
76 position,
77 descending: key.descending,
78 nulls_first: Some(key.nulls_first.unwrap_or(key.descending)),
79 nullable: true,
80 })
81 })
82 .collect::<Option<Vec<_>>>()
83 .unwrap_or_default();
84 Self {
85 child,
86 keys,
87 evaluator,
88 keep,
89 work_mem_bytes,
90 spill_directory: None,
91 schema,
92 input_slots,
93 run_schema,
94 ordering,
95 output: None,
96 initial_run_count: 0,
97 merge_pass_count: 0,
98 }
99 }
100
101 pub fn with_spill_directory(mut self, directory: impl Into<PathBuf>) -> Self {
103 self.spill_directory = Some(directory.into());
104 self
105 }
106
107 pub fn initial_run_count(&self) -> usize {
108 self.initial_run_count
109 }
110
111 pub fn merge_pass_count(&self) -> usize {
112 self.merge_pass_count
113 }
114
115 fn create_run_buffer(&self) -> SpillBuffer {
116 self.spill_directory.as_ref().map_or_else(
117 || SpillBuffer::new(self.work_mem_bytes),
118 |directory| SpillBuffer::new_in(self.work_mem_bytes, directory),
119 )
120 }
121
122 fn build_initial_runs(&mut self) -> ExecResult<Vec<SortedRun>> {
123 let mut sequence = 0_u64;
124 let mut pending = Vec::new();
125 let mut pending_size = EncodedBatchSizer::new(&self.run_schema)?;
126 let mut runs = Vec::new();
127
128 while let Some(batch) = self.child.next()? {
129 for row in batch.rows {
130 let mut key_values = Vec::with_capacity(self.keys.len());
131 for key in &self.keys {
132 key_values.push(self.evaluator.evaluate_physical(
133 &key.expr,
134 &batch.schema,
135 &row,
136 )?);
137 }
138 let record_sequence = sequence;
139 key_values.push(Value::Bytes(record_sequence.to_be_bytes().to_vec()));
140 let record = DecoratedRow {
141 row: row
142 .project_slots(&self.input_slots)
143 .append_values(key_values),
144 sequence,
145 };
146 sequence = sequence.checked_add(1).ok_or_else(|| {
147 ExecError::Other("external sort input sequence overflow".into())
148 })?;
149 let mut candidate_size = pending_size;
150 candidate_size.append(&record.row)?;
151 let would_exceed = candidate_size.bytes() > self.work_mem_bytes;
152
153 if would_exceed && !pending.is_empty() {
154 if let Some(run) = self.finish_run(std::mem::take(&mut pending), true)? {
155 runs.push(run);
156 }
157 pending_size = EncodedBatchSizer::new(&self.run_schema)?;
158 candidate_size = pending_size;
159 candidate_size.append(&record.row)?;
160 }
161
162 pending.push(record);
163 pending_size = candidate_size;
164
165 if pending_size.bytes() > self.work_mem_bytes {
168 if let Some(run) = self.finish_run(std::mem::take(&mut pending), true)? {
169 runs.push(run);
170 }
171 pending_size = EncodedBatchSizer::new(&self.run_schema)?;
172 }
173 }
174 }
175
176 let force_spill = !runs.is_empty();
177 if let Some(run) = self.finish_run(pending, force_spill)? {
178 runs.push(run);
179 }
180 Ok(runs)
181 }
182
183 fn finish_run(
184 &self,
185 records: Vec<DecoratedRow>,
186 force_spill: bool,
187 ) -> ExecResult<Option<SortedRun>> {
188 let mut records = records;
189 records.sort_unstable_by(|left, right| {
190 compare_records(
191 &self.keys,
192 &self.run_schema,
193 self.input_slots.len(),
194 left,
195 right,
196 )
197 });
198 if let Some(keep) = self.keep {
199 records.truncate(keep);
200 }
201 if records.is_empty() {
202 return Ok(None);
203 }
204
205 let mut buffer = self.create_run_buffer();
206 let mut writer = RunBatchWriter::new(self.run_schema.clone())?;
207 for record in records {
208 writer.push(&mut buffer, record.row)?;
209 }
210 writer.finish(&mut buffer)?;
211 if force_spill {
212 buffer.spill_pending()?;
213 }
214 Ok(Some(SortedRun { buffer }))
215 }
216
217 fn collapse_runs(&mut self, mut runs: Vec<SortedRun>) -> ExecResult<Option<SortedRun>> {
218 while runs.len() > 1 {
219 self.merge_pass_count = self.merge_pass_count.checked_add(1).ok_or_else(|| {
220 ExecError::Other("external sort merge-pass count overflow".into())
221 })?;
222 let mut inputs = runs.into_iter();
223 let mut merged = Vec::new();
224 loop {
225 let group: Vec<_> = inputs.by_ref().take(EXTERNAL_SORT_MERGE_FAN_IN).collect();
226 if group.is_empty() {
227 break;
228 }
229 merged.push(merge_group(
230 group,
231 &self.keys,
232 self.keep,
233 self.create_run_buffer(),
234 &self.run_schema,
235 self.input_slots.len(),
236 )?);
237 }
238 runs = merged;
239 }
240 Ok(runs.pop())
241 }
242}
243
244impl PhysicalOperator for ExternalSort<'_> {
245 fn row_schema(&self) -> &RowSchema {
246 &self.schema
247 }
248
249 fn output_ordering(&self) -> &[PhysicalOrder] {
250 &self.ordering
251 }
252
253 fn open(&mut self) -> ExecResult<()> {
254 self.output = None;
255 self.initial_run_count = 0;
256 self.merge_pass_count = 0;
257 self.child.open()?;
258
259 let runs = self.build_initial_runs()?;
260 self.initial_run_count = runs.len();
261 let Some(mut final_run) = self.collapse_runs(runs)? else {
262 return Ok(());
263 };
264 self.output = Some(final_run.buffer.drain()?);
265 Ok(())
266 }
267
268 fn next(&mut self) -> ExecResult<Option<Batch>> {
269 let Some(output) = self.output.as_mut() else {
270 return Ok(None);
271 };
272 loop {
273 let Some(batch) = output.next().transpose()? else {
274 return Ok(None);
275 };
276 validate_run_batch(&batch, &self.run_schema)?;
277 if batch.rows.is_empty() {
278 continue;
279 }
280 let rows = batch
281 .rows
282 .into_iter()
283 .map(|row| row.into_prefix(self.input_slots.len()))
284 .collect();
285 return Ok(Some(Batch::from_physical_rows(self.schema.clone(), rows)));
286 }
287 }
288
289 fn close(&mut self) -> ExecResult<()> {
290 self.output = None;
291 self.child.close()
292 }
293}
294
295struct SortedRun {
296 buffer: SpillBuffer,
297}
298
299struct DecoratedRow {
300 row: PhysicalRow,
301 sequence: u64,
302}
303
304fn decode_record(
305 schema: &RowSchema,
306 row: PhysicalRow,
307 source_width: usize,
308 expected_key_count: usize,
309) -> ExecResult<DecoratedRow> {
310 let sequence_position = source_width
311 .checked_add(expected_key_count)
312 .ok_or_else(|| ExecError::Other("external sort run width overflow".into()))?;
313 if schema.len() != sequence_position.saturating_add(1) {
314 return Err(ExecError::Other(format!(
315 "invalid external sort run key count: expected {expected_key_count}"
316 )));
317 }
318 let view = schema.view(&row);
319 for index in source_width..sequence_position {
320 if view.value_at(index).is_none() {
321 return Err(ExecError::Other(format!(
322 "external sort run is missing key {}",
323 index - source_width
324 )));
325 }
326 }
327 let sequence = match view.value_at(sequence_position) {
328 Some(Value::Bytes(bytes)) if bytes.len() == std::mem::size_of::<u64>() => {
329 let bytes: [u8; 8] = bytes
330 .as_slice()
331 .try_into()
332 .map_err(|_| ExecError::Other("invalid external sort run sequence width".into()))?;
333 u64::from_be_bytes(bytes)
334 }
335 _ => {
336 return Err(ExecError::Other(
337 "invalid external sort run sequence".into(),
338 ))
339 }
340 };
341 Ok(DecoratedRow { row, sequence })
342}
343
344fn validate_run_batch(batch: &Batch, expected: &RowSchema) -> ExecResult<()> {
345 if &batch.schema == expected {
346 Ok(())
347 } else {
348 Err(ExecError::Other(format!(
349 "invalid external sort run schema: expected {:?}, got {:?}",
350 expected.columns(),
351 batch.schema.columns()
352 )))
353 }
354}
355
356fn compare_records(
357 keys: &[SortKey],
358 schema: &RowSchema,
359 source_width: usize,
360 left: &DecoratedRow,
361 right: &DecoratedRow,
362) -> Ordering {
363 let left_view = schema.view(&left.row);
364 let right_view = schema.view(&right.row);
365 compare_sort_key_values_by(keys, |index| {
366 (
367 left_view
368 .value_at(source_width + index)
369 .expect("validated external sort run key"),
370 right_view
371 .value_at(source_width + index)
372 .expect("validated external sort run key"),
373 )
374 })
375 .then_with(|| left.sequence.cmp(&right.sequence))
376}
377
378struct RunBatchWriter {
379 schema: RowSchema,
380 pending: Vec<PhysicalRow>,
381 pending_size: EncodedBatchSizer,
382}
383
384impl RunBatchWriter {
385 fn new(schema: RowSchema) -> ExecResult<Self> {
386 let pending_size = EncodedBatchSizer::new(&schema)?;
387 Ok(Self {
388 schema,
389 pending: Vec::new(),
390 pending_size,
391 })
392 }
393
394 fn push(&mut self, output: &mut SpillBuffer, record: PhysicalRow) -> ExecResult<()> {
395 let mut candidate_size = self.pending_size;
396 candidate_size.append(&record)?;
397 let exceeds_budget = candidate_size.bytes() > output.budget_bytes();
398 if exceeds_budget && !self.pending.is_empty() {
399 self.flush(output)?;
400 candidate_size = self.pending_size;
401 candidate_size.append(&record)?;
402 }
403 self.pending.push(record);
404 self.pending_size = candidate_size;
405 if self.pending_size.bytes() > output.budget_bytes() {
406 self.flush(output)?;
407 }
408 Ok(())
409 }
410
411 fn finish(mut self, output: &mut SpillBuffer) -> ExecResult<()> {
412 self.flush(output)
413 }
414
415 fn flush(&mut self, output: &mut SpillBuffer) -> ExecResult<()> {
416 if self.pending.is_empty() {
417 return Ok(());
418 }
419 output.push(Batch::from_physical_rows(
420 self.schema.clone(),
421 std::mem::take(&mut self.pending),
422 ))?;
423 self.pending_size = EncodedBatchSizer::new(&self.schema)?;
424 Ok(())
425 }
426}
427
428struct MergeCursor {
429 batches: SpillDrain,
430 rows: std::vec::IntoIter<PhysicalRow>,
431 schema: RowSchema,
432 source_width: usize,
433 key_count: usize,
434}
435
436impl MergeCursor {
437 fn new(
438 mut buffer: SpillBuffer,
439 schema: RowSchema,
440 source_width: usize,
441 key_count: usize,
442 ) -> ExecResult<Self> {
443 Ok(Self {
444 batches: buffer.drain()?,
445 rows: Vec::new().into_iter(),
446 schema,
447 source_width,
448 key_count,
449 })
450 }
451
452 fn next_record(&mut self) -> ExecResult<Option<DecoratedRow>> {
453 loop {
454 if let Some(row) = self.rows.next() {
455 return decode_record(&self.schema, row, self.source_width, self.key_count)
456 .map(Some);
457 }
458 let Some(batch) = self.batches.next().transpose()? else {
459 return Ok(None);
460 };
461 validate_run_batch(&batch, &self.schema)?;
462 self.rows = batch.rows.into_iter();
463 }
464 }
465}
466
467struct HeapItem {
468 record: DecoratedRow,
469 cursor: usize,
470}
471
472fn merge_group(
473 runs: Vec<SortedRun>,
474 keys: &[SortKey],
475 keep: Option<usize>,
476 mut output: SpillBuffer,
477 run_schema: &RowSchema,
478 source_width: usize,
479) -> ExecResult<SortedRun> {
480 debug_assert!(!runs.is_empty());
481 debug_assert!(runs.len() <= EXTERNAL_SORT_MERGE_FAN_IN);
482 let mut cursors = Vec::with_capacity(runs.len());
483 let mut heap = Vec::with_capacity(runs.len());
484
485 for run in runs {
486 cursors.push(MergeCursor::new(
487 run.buffer,
488 run_schema.clone(),
489 source_width,
490 keys.len(),
491 )?);
492 let cursor = cursors.len() - 1;
493 if let Some(record) = cursors[cursor].next_record()? {
494 heap_push(
495 &mut heap,
496 HeapItem { record, cursor },
497 keys,
498 run_schema,
499 source_width,
500 );
501 }
502 }
503
504 let mut writer = RunBatchWriter::new(run_schema.clone())?;
505 let mut emitted = 0_usize;
506 while !heap.is_empty() && keep.is_none_or(|keep| emitted < keep) {
507 let item = heap_pop(&mut heap, keys, run_schema, source_width)
508 .ok_or_else(|| ExecError::Other("external sort merge heap became empty".into()))?;
509 let cursor = item.cursor;
510 writer.push(&mut output, item.record.row)?;
511 emitted = emitted
512 .checked_add(1)
513 .ok_or_else(|| ExecError::Other("external sort emitted-row count overflow".into()))?;
514 if let Some(record) = cursors[cursor].next_record()? {
515 heap_push(
516 &mut heap,
517 HeapItem { record, cursor },
518 keys,
519 run_schema,
520 source_width,
521 );
522 }
523 }
524 writer.finish(&mut output)?;
525 output.spill_pending()?;
526 Ok(SortedRun { buffer: output })
527}
528
529fn compare_heap_items(
530 keys: &[SortKey],
531 schema: &RowSchema,
532 source_width: usize,
533 left: &HeapItem,
534 right: &HeapItem,
535) -> Ordering {
536 compare_records(keys, schema, source_width, &left.record, &right.record)
537}
538
539fn heap_push(
540 heap: &mut Vec<HeapItem>,
541 item: HeapItem,
542 keys: &[SortKey],
543 schema: &RowSchema,
544 source_width: usize,
545) {
546 heap.push(item);
547 let mut child = heap.len() - 1;
548 while child > 0 {
549 let parent = (child - 1) / 2;
550 if compare_heap_items(keys, schema, source_width, &heap[child], &heap[parent])
551 != Ordering::Less
552 {
553 break;
554 }
555 heap.swap(child, parent);
556 child = parent;
557 }
558}
559
560fn heap_pop(
561 heap: &mut Vec<HeapItem>,
562 keys: &[SortKey],
563 schema: &RowSchema,
564 source_width: usize,
565) -> Option<HeapItem> {
566 if heap.is_empty() {
567 return None;
568 }
569 let smallest = heap.swap_remove(0);
570 let mut parent = 0;
571 loop {
572 let left = parent * 2 + 1;
573 if left >= heap.len() {
574 break;
575 }
576 let right = left + 1;
577 let child = if right < heap.len()
578 && compare_heap_items(keys, schema, source_width, &heap[right], &heap[left])
579 == Ordering::Less
580 {
581 right
582 } else {
583 left
584 };
585 if compare_heap_items(keys, schema, source_width, &heap[child], &heap[parent])
586 != Ordering::Less
587 {
588 break;
589 }
590 heap.swap(parent, child);
591 parent = child;
592 }
593 Some(smallest)
594}
595
596#[cfg(test)]
597mod tests {
598 use std::collections::BTreeMap;
599 use std::io::Write as _;
600 use std::sync::Arc;
601
602 use super::*;
603 use crate::physical::{run_to_batches, run_to_rows};
604 use crate::scalar::ScalarExpr;
605 use crate::scan::TableScan;
606 use uqa_sql::expr::RowLookup as _;
607 use uqa_sql::ResultRow;
608
609 struct Columns;
610
611 struct PhysicalRowsScan {
612 schema: RowSchema,
613 rows: Option<Vec<PhysicalRow>>,
614 }
615
616 impl PhysicalOperator for PhysicalRowsScan {
617 fn row_schema(&self) -> &RowSchema {
618 &self.schema
619 }
620
621 fn open(&mut self) -> ExecResult<()> {
622 Ok(())
623 }
624
625 fn next(&mut self) -> ExecResult<Option<Batch>> {
626 Ok(self
627 .rows
628 .take()
629 .map(|rows| Batch::from_physical_rows(self.schema.clone(), rows)))
630 }
631
632 fn close(&mut self) -> ExecResult<()> {
633 Ok(())
634 }
635 }
636
637 impl crate::relational::ExpressionEvaluator for Columns {
638 fn evaluate(
639 &self,
640 expression: &ScalarExpr,
641 row: &dyn uqa_sql::expr::RowLookup,
642 ) -> ExecResult<Value> {
643 match expression {
644 ScalarExpr::Column(name) => Ok(row.column(name).cloned().unwrap_or(Value::Null)),
645 _ => Err(ExecError::Other(
646 "test evaluator only supports columns".into(),
647 )),
648 }
649 }
650 }
651
652 fn row(key: i64, input: i64) -> ResultRow {
653 BTreeMap::from([
654 ("key".into(), Value::Int(key)),
655 ("input".into(), Value::Int(input)),
656 ])
657 }
658
659 fn sort(rows: Vec<ResultRow>, budget: usize, keep: Option<usize>) -> ExternalSort<'static> {
660 ExternalSort::new(
661 Box::new(TableScan::from_rows(
662 vec!["key".into(), "input".into()],
663 rows,
664 )),
665 vec![SortKey {
666 expr: ScalarExpr::Column("key".into()),
667 descending: false,
668 nulls_first: None,
669 }],
670 Arc::new(Columns),
671 keep,
672 budget,
673 )
674 }
675
676 fn int_column(rows: &[ResultRow], column: &str) -> Vec<i64> {
677 rows.iter()
678 .map(|row| match row.get(column) {
679 Some(Value::Int(value)) => *value,
680 value => panic!("unexpected {column} value: {value:?}"),
681 })
682 .collect()
683 }
684
685 #[test]
686 fn tiny_budget_builds_many_runs_and_multi_pass_merge() {
687 let rows = (0..(EXTERNAL_SORT_MERGE_FAN_IN as i64 * 2 + 5))
688 .rev()
689 .map(|value| row(value, value))
690 .collect();
691 let mut operator = sort(rows, 1, None);
692 operator.open().unwrap();
693 assert!(operator.initial_run_count() > EXTERNAL_SORT_MERGE_FAN_IN);
694 assert!(operator.merge_pass_count() >= 2);
695 let mut output = Vec::new();
696 while let Some(batch) = operator.next().unwrap() {
697 output.extend(batch.into_result_rows());
698 }
699 operator.close().unwrap();
700 assert_eq!(
701 int_column(&output, "key"),
702 (0..(EXTERNAL_SORT_MERGE_FAN_IN as i64 * 2 + 5)).collect::<Vec<_>>()
703 );
704 }
705
706 #[test]
707 fn equal_keys_keep_original_input_order_across_runs() {
708 let rows = (0..40).map(|input| row(7, input)).collect();
709 let mut operator = sort(rows, 1, None);
710 let (_, output) = run_to_rows(&mut operator).unwrap();
711 assert_eq!(int_column(&output, "input"), (0..40).collect::<Vec<_>>());
712 }
713
714 #[test]
715 fn keep_is_global_top_k_across_multiple_runs() {
716 let rows = (0..100).rev().map(|value| row(value, value)).collect();
717 let mut operator = sort(rows, 1, Some(7));
718 let (_, output) = run_to_rows(&mut operator).unwrap();
719 assert_eq!(int_column(&output, "key"), (0..7).collect::<Vec<_>>());
720 }
721
722 #[test]
723 fn spill_creation_error_is_propagated() {
724 let not_a_directory = tempfile::NamedTempFile::new().unwrap();
725 let mut operator = sort(vec![row(1, 0)], 0, None)
726 .with_spill_directory(not_a_directory.path().to_path_buf());
727 let error = operator.open().unwrap_err();
728 assert!(error.to_string().contains("failed to create spill file"));
729 }
730
731 #[test]
732 fn large_budget_single_run_does_not_create_a_spill_file() {
733 let not_a_directory = tempfile::NamedTempFile::new().unwrap();
734 let rows = (0..20).rev().map(|value| row(value, value)).collect();
735 let mut operator =
736 sort(rows, 1_000_000, None).with_spill_directory(not_a_directory.path().to_path_buf());
737 let (_, output) = run_to_rows(&mut operator).unwrap();
738 assert_eq!(int_column(&output, "key"), (0..20).collect::<Vec<_>>());
739 assert_eq!(operator.initial_run_count(), 1);
740 assert_eq!(operator.merge_pass_count(), 0);
741 }
742
743 #[test]
744 fn spill_batch_schema_overhead_is_amortized_across_sort_rows() {
745 let rows = (0..20_000).rev().map(|value| row(value, value)).collect();
746 let mut operator = sort(rows, 4 * 1024 * 1024, None);
747
748 operator.open().unwrap();
749 assert_eq!(operator.initial_run_count(), 1);
750 assert_eq!(operator.merge_pass_count(), 0);
751 operator.close().unwrap();
752 }
753
754 #[test]
755 fn mixed_lock_origins_add_metadata_for_origin_free_rows_to_the_batch_budget() {
756 let schema = RowSchema::new(vec!["value".into()]);
757 let plain = PhysicalRow::from_values(vec![Value::Int(1)]);
758 let locked = PhysicalRow::from_values(vec![Value::Int(2)])
759 .with_lock_origin(crate::RowLockOrigin::new("source", "public.source", 2));
760 let overhead = EncodedBatchSizer::new(&schema).unwrap().bytes();
761 let mut plain_size = EncodedBatchSizer::new(&schema).unwrap();
762 plain_size.append(&plain).unwrap();
763 let plain_record_bytes = plain_size.bytes() - overhead;
764 let mut locked_size = EncodedBatchSizer::new(&schema).unwrap();
765 locked_size.append(&locked).unwrap();
766 let locked_record_bytes = locked_size.bytes() - overhead;
767 let mut size = EncodedBatchSizer::new(&schema).unwrap();
768 size.append(&plain).unwrap();
769 let without_retroactive_metadata = overhead + plain_record_bytes + locked_record_bytes;
770 size.append(&locked).unwrap();
771 let batch = Batch::from_physical_rows(schema, vec![plain, locked]);
772
773 assert_eq!(size.bytes(), SpillBuffer::encoded_size(&batch).unwrap());
774 assert_eq!(size.bytes(), without_retroactive_metadata + 8);
775 }
776
777 #[test]
778 fn corrupt_run_read_error_is_propagated() {
779 let keys = vec![SortKey {
780 expr: ScalarExpr::Column("key".into()),
781 descending: false,
782 nulls_first: None,
783 }];
784 let schema = run_schema(2, 1);
785 let record = PhysicalRow::from_values(vec![
786 Value::Int(1),
787 Value::Int(0),
788 Value::Int(1),
789 Value::Bytes(0_u64.to_be_bytes().to_vec()),
790 ]);
791 let mut buffer = SpillBuffer::new(0);
792 buffer
793 .push(Batch::from_physical_rows(schema.clone(), vec![record]))
794 .unwrap();
795 let path = buffer.spill_path().unwrap().to_path_buf();
796 let mut corrupt = std::fs::OpenOptions::new().append(true).open(path).unwrap();
797 corrupt.write_all(&1_u64.to_le_bytes()).unwrap();
798 corrupt.write_all(&[0xff]).unwrap();
799 corrupt.flush().unwrap();
800
801 let result = merge_group(
802 vec![SortedRun { buffer }],
803 &keys,
804 None,
805 SpillBuffer::new(0),
806 &schema,
807 2,
808 );
809 let error = match result {
810 Ok(_) => panic!("corrupt run unexpectedly merged"),
811 Err(error) => error,
812 };
813 assert!(error
814 .to_string()
815 .contains("truncated schema physical width"));
816 }
817
818 #[test]
819 fn corrupt_run_key_width_is_reported_before_comparison() {
820 let schema = run_schema(2, 0);
821 let record = PhysicalRow::from_values(vec![
822 Value::Int(1),
823 Value::Int(0),
824 Value::Bytes(0_u64.to_be_bytes().to_vec()),
825 ]);
826 let error = match decode_record(&schema, record, 2, 1) {
827 Ok(_) => panic!("corrupt run record unexpectedly decoded"),
828 Err(error) => error,
829 };
830 assert!(error
831 .to_string()
832 .contains("invalid external sort run key count"));
833 }
834
835 #[test]
836 fn forced_spill_preserves_hidden_alias_and_public_column_slots() {
837 let base = RowSchema::new(vec!["value".into(), "key".into()]);
838 let aliased = RowSchema::with_identity_aliases(
839 &base,
840 &[(crate::ColumnIdentity::qualified("source", "value"), 0)],
841 );
842 let schema = RowSchema::append(&aliased, &["value".into()]);
843 let rows = vec![
844 PhysicalRow::from_values(vec![Value::Str("source-b".into()), Value::Int(2)])
845 .append_values(vec![Value::Str("projected-b".into())]),
846 PhysicalRow::from_values(vec![Value::Str("source-a".into()), Value::Int(1)])
847 .append_values(vec![Value::Str("projected-a".into())]),
848 ];
849 let scan = PhysicalRowsScan {
850 schema,
851 rows: Some(rows),
852 };
853 let mut operator = ExternalSort::new(
854 Box::new(scan),
855 vec![SortKey {
856 expr: ScalarExpr::Column("key".into()),
857 descending: false,
858 nulls_first: None,
859 }],
860 Arc::new(Columns),
861 None,
862 1,
863 );
864
865 let batches = run_to_batches(&mut operator).unwrap();
866 assert!(operator.initial_run_count() > 1);
867 let values = batches
868 .iter()
869 .flat_map(|batch| {
870 batch.rows.iter().map(|row| {
871 let view = batch.schema.view(row);
872 (
873 view.get("value").cloned(),
874 view.qualified_column("source", "value").cloned(),
875 )
876 })
877 })
878 .collect::<Vec<_>>();
879 assert_eq!(
880 values,
881 vec![
882 (
883 Some(Value::Str("projected-a".into())),
884 Some(Value::Str("source-a".into()))
885 ),
886 (
887 Some(Value::Str("projected-b".into())),
888 Some(Value::Str("source-b".into()))
889 ),
890 ]
891 );
892 }
893
894 #[test]
895 fn forced_spill_preserves_duplicate_logical_columns_positionally() {
896 let schema = RowSchema::new(vec!["value".into(), "value".into(), "key".into()]);
897 let scan = PhysicalRowsScan {
898 schema,
899 rows: Some(vec![
900 PhysicalRow::from_values(vec![
901 Value::Str("left-b".into()),
902 Value::Str("right-b".into()),
903 Value::Int(2),
904 ]),
905 PhysicalRow::from_values(vec![
906 Value::Str("left-a".into()),
907 Value::Str("right-a".into()),
908 Value::Int(1),
909 ]),
910 ]),
911 };
912 let mut operator = ExternalSort::new(
913 Box::new(scan),
914 vec![SortKey {
915 expr: ScalarExpr::Column("key".into()),
916 descending: false,
917 nulls_first: None,
918 }],
919 Arc::new(Columns),
920 None,
921 1,
922 );
923
924 let batches = run_to_batches(&mut operator).unwrap();
925 assert!(operator.initial_run_count() > 1);
926 assert_eq!(operator.schema(), ["value", "value", "key"]);
927 let values = batches
928 .iter()
929 .flat_map(|batch| {
930 batch.rows.iter().map(|row| {
931 let view = batch.schema.view(row);
932 (view.value_at(0).cloned(), view.value_at(1).cloned())
933 })
934 })
935 .collect::<Vec<_>>();
936 assert_eq!(
937 values,
938 vec![
939 (
940 Some(Value::Str("left-a".into())),
941 Some(Value::Str("right-a".into()))
942 ),
943 (
944 Some(Value::Str("left-b".into())),
945 Some(Value::Str("right-b".into()))
946 ),
947 ]
948 );
949 }
950
951 #[test]
952 fn empty_input_and_zero_keep_are_empty() {
953 let mut empty = sort(Vec::new(), 1, None);
954 assert!(run_to_rows(&mut empty).unwrap().1.is_empty());
955
956 let rows = (0..10).map(|value| row(value, value)).collect();
957 let mut zero = sort(rows, 1, Some(0));
958 assert!(run_to_rows(&mut zero).unwrap().1.is_empty());
959 }
960}