Skip to main content

uqa_execution/
external_sort.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Byte-bounded external merge sort.
8//!
9//! Input is divided into stable sorted runs whose exact spill encoding stays
10//! within `work_mem_bytes` (except for one indivisible oversized row). This is
11//! a hard bound on retained encoded row/run bytes, not Rust allocator resident
12//! bytes. Runs are written through [`crate::spill::SpillBuffer`], then merged
13//! with a fixed fan-in. Merge reader/heap overhead is therefore constant: at
14//! most [`EXTERNAL_SORT_MERGE_FAN_IN`] decoded rows plus one byte-bounded output
15//! buffer, independent of the total run count.
16
17use 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
29/// Maximum number of input runs opened by one merge operation.
30pub 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
39/// Physical external sort with stable SQL ordering and optional global top-K.
40pub 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    /// Place sort runs in a caller-selected temporary-data directory.
96    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                // One row is indivisible. Flush it immediately when it alone is
160                // larger than work_mem so no additional row joins it in memory.
161                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}