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