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
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
45/// Physical external sort with stable SQL ordering and optional global top-K.
46pub 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    /// Place sort runs in a caller-selected temporary-data directory.
102    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                // One row is indivisible. Flush it immediately when it alone is
166                // larger than work_mem so no additional row joins it in memory.
167                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}