#![allow(dead_code)]
use alloc::boxed::Box;
use alloc::vec::Vec;
use spg_storage::{Row, TableSchema, Value};
use crate::EngineError;
use crate::orderby::{OrderKey, cmp_multi_key_in};
use crate::tempstore::TempRun;
type Head = Option<(Vec<OrderKey>, Row<'static>)>;
const NO_RUN: usize = usize::MAX;
fn head_cmp(heads: &[Head], descs: &[bool], a: usize, b: usize) -> core::cmp::Ordering {
use core::cmp::Ordering;
match (a == NO_RUN, b == NO_RUN) {
(true, true) => return Ordering::Equal,
(true, false) => return Ordering::Greater,
(false, true) => return Ordering::Less,
(false, false) => {}
}
match (&heads[a], &heads[b]) {
(None, None) => Ordering::Equal,
(None, Some(_)) => Ordering::Greater,
(Some(_), None) => Ordering::Less,
(Some((ka, _)), Some((kb, _))) => cmp_multi_key_in(ka, kb, descs, &[]),
}
}
struct HeadHeap {
heap: Vec<usize>,
}
impl HeadHeap {
fn build(k: usize, heads: &[Head], descs: &[bool]) -> Self {
let mut h = Self {
heap: (0..k).filter(|&i| heads[i].is_some()).collect(),
};
for start in (0..h.heap.len() / 2).rev() {
h.sift_down(start, heads, descs);
}
h
}
fn peek(&self) -> Option<usize> {
self.heap.first().copied()
}
fn settle_root(&mut self, heads: &[Head], descs: &[bool]) {
let Some(&top) = self.heap.first() else {
return;
};
if heads[top].is_none() {
let last = self.heap.pop();
if let Some(last) = last
&& !self.heap.is_empty()
{
self.heap[0] = last;
}
}
if !self.heap.is_empty() {
self.sift_down(0, heads, descs);
}
}
fn sift_down(&mut self, mut node: usize, heads: &[Head], descs: &[bool]) {
let n = self.heap.len();
loop {
let (l, r) = (2 * node + 1, 2 * node + 2);
let mut best = node;
if l < n
&& head_cmp(heads, descs, self.heap[l], self.heap[best])
== core::cmp::Ordering::Less
{
best = l;
}
if r < n
&& head_cmp(heads, descs, self.heap[r], self.heap[best])
== core::cmp::Ordering::Less
{
best = r;
}
if best == node {
return;
}
self.heap.swap(node, best);
node = best;
}
}
}
fn read_exact(run: &mut dyn TempRun, n: usize) -> Result<Option<Vec<u8>>, EngineError> {
let mut buf = alloc::vec![0u8; n];
let mut filled = 0;
while filled < n {
let got = run
.read(&mut buf[filled..])
.map_err(|e| EngineError::Internal(alloc::format!("temp run read: {e:?}")))?;
if got == 0 {
if filled == 0 {
return Ok(None);
}
return Err(EngineError::Internal(alloc::string::String::from(
"temp run ended mid-record",
)));
}
filled += got;
}
Ok(Some(buf))
}
pub(crate) struct ExternalSorter<'a> {
factory: Option<crate::TempRunFactory>,
budget_bytes: usize,
record_schema: TableSchema,
descs: &'a [bool],
arena: Vec<u8>,
ends: Vec<u32>,
keys: Vec<OrderKey>,
key_stride: usize,
runs: Vec<Box<dyn TempRun>>,
stats: Option<&'a crate::tempstore::SpillStats>,
needed: &'a [bool],
}
impl<'a> ExternalSorter<'a> {
pub(crate) fn new(
factory: Option<crate::TempRunFactory>,
budget_bytes: usize,
record_cols: Vec<spg_storage::ColumnSchema>,
descs: &'a [bool],
) -> Self {
Self {
factory,
budget_bytes,
record_schema: TableSchema::new("spg_sort_run", record_cols),
descs,
arena: Vec::new(),
ends: Vec::new(),
keys: Vec::new(),
key_stride: 0,
runs: Vec::new(),
stats: None,
needed: &[],
}
}
pub(crate) fn with_pruned(mut self, needed: &'a [bool]) -> Self {
debug_assert!(
needed.is_empty() || needed.len() == self.record_schema.columns.len(),
"prune mask must match the record arity"
);
self.needed = needed;
self
}
pub(crate) fn spilled(&self) -> bool {
!self.runs.is_empty()
}
pub(crate) fn push(
&mut self,
keys: &mut Vec<OrderKey>,
record: &Row<'_>,
) -> Result<(), EngineError> {
if self.ends.is_empty() {
self.key_stride = keys.len();
}
debug_assert_eq!(
keys.len(),
self.key_stride,
"every row of one sort carries the same number of keys"
);
self.keys.append(keys);
spg_storage::encode_row_body_dense_masked_into(
record,
&self.record_schema,
self.needed,
&mut self.arena,
);
let end = u32::try_from(self.arena.len()).map_err(|_| {
EngineError::Internal(alloc::string::String::from("sort batch larger than 4 GiB"))
})?;
self.ends.push(end);
if self.factory.is_some() && self.batch_bytes() >= self.budget_bytes {
self.spill_batch()?;
}
Ok(())
}
pub(crate) fn with_stats(mut self, stats: &'a crate::tempstore::SpillStats) -> Self {
self.stats = Some(stats);
self
}
fn batch_bytes(&self) -> usize {
self.arena.len()
+ self.keys.len() * core::mem::size_of::<OrderKey>()
+ self.ends.len() * core::mem::size_of::<u32>()
}
fn span(&self, i: usize) -> (usize, usize) {
let start = if i == 0 { 0 } else { self.ends[i - 1] as usize };
(start, self.ends[i] as usize)
}
fn sorted_order(&self) -> Vec<u32> {
let mut order: Vec<u32> = (0..self.ends.len() as u32).collect();
if self.key_stride == 0 {
return order;
}
let (keys, stride, descs) = (&self.keys, self.key_stride, self.descs);
if stride == 1
&& let Some(mut inline) = Self::inline_int_keys(keys)
{
if descs.first().copied().unwrap_or(false) {
inline.sort_by_key(|p| core::cmp::Reverse(p.0));
} else {
inline.sort_by_key(|p| p.0);
}
for (slot, (_, i)) in order.iter_mut().zip(inline) {
*slot = i;
}
return order;
}
order.sort_by(|&a, &b| {
let (a, b) = (a as usize * stride, b as usize * stride);
cmp_multi_key_in(&keys[a..a + stride], &keys[b..b + stride], descs, &[])
});
order
}
fn inline_int_keys(keys: &[OrderKey]) -> Option<Vec<(i128, u32)>> {
let mut out: Vec<(i128, u32)> = Vec::with_capacity(keys.len());
for (i, k) in keys.iter().enumerate() {
out.push((crate::orderby::inline_int_key(k)?, i as u32));
}
Some(out)
}
fn clear_batch(&mut self) {
self.arena.clear();
self.ends.clear();
self.keys.clear();
}
fn spill_batch(&mut self) -> Result<(), EngineError> {
if self.ends.is_empty() {
return Ok(());
}
let Some(factory) = self.factory else {
return Ok(());
};
let order = self.sorted_order();
let mut run = factory()
.map_err(|e| EngineError::Internal(alloc::format!("temp run create: {e:?}")))?;
for i in order {
let (start, end) = self.span(i as usize);
let len = u32::try_from(end - start).map_err(|_| {
EngineError::Internal(alloc::string::String::from("row too large to spill"))
})?;
run.append(&len.to_le_bytes())
.map_err(|e| EngineError::Internal(alloc::format!("temp run append: {e:?}")))?;
run.append(&self.arena[start..end])
.map_err(|e| EngineError::Internal(alloc::format!("temp run append: {e:?}")))?;
}
run.seal()
.map_err(|e| EngineError::Internal(alloc::format!("temp run seal: {e:?}")))?;
if let Some(stats) = self.stats {
use core::sync::atomic::Ordering;
stats.files.fetch_add(1, Ordering::Relaxed);
stats
.bytes
.fetch_add(run.bytes_written(), Ordering::Relaxed);
}
self.runs.push(run);
self.clear_batch();
Ok(())
}
pub(crate) fn finish_each<K, P, E>(
mut self,
keys_of: K,
project: P,
mut emit: E,
) -> Result<usize, EngineError>
where
K: Fn(&Row<'static>) -> Result<Vec<OrderKey>, EngineError>,
P: Fn(&Row<'static>, &mut Vec<Value<'static>>) -> Result<(), EngineError>,
E: FnMut(&[Value<'static>]) -> Result<(), EngineError>,
{
let mut emitted = 0usize;
let needed = self.needed;
let mut scratch: Vec<Value<'static>> = Vec::new();
if self.runs.is_empty() {
for i in self.sorted_order() {
let (start, end) = self.span(i as usize);
let (row, _) = spg_storage::decode_row_body_dense_pruned(
&self.arena[start..end],
&self.record_schema,
spg_storage::CURRENT_ROW_CODEC_VERSION,
needed,
)
.map_err(|e| EngineError::Internal(alloc::format!("sort batch decode: {e:?}")))?;
scratch.clear();
project(&row, &mut scratch)?;
emit(&scratch)?;
emitted += 1;
}
return Ok(emitted);
}
self.spill_batch()?;
let mut runs: Vec<Box<dyn TempRun>> = core::mem::take(&mut self.runs);
let mut heads: Vec<Head> = Vec::with_capacity(runs.len());
for run in &mut runs {
heads.push(Self::next_row(
&mut **run,
&self.record_schema,
&keys_of,
needed,
)?);
}
let mut heap = HeadHeap::build(heads.len(), &heads, self.descs);
while let Some(w) = heap.peek() {
let Some((_, record)) = heads[w].take() else {
break;
};
scratch.clear();
project(&record, &mut scratch)?;
emit(&scratch)?;
emitted += 1;
heads[w] = Self::next_row(&mut *runs[w], &self.record_schema, &keys_of, needed)?;
heap.settle_root(&heads, self.descs);
}
Ok(emitted)
}
pub(crate) fn finish<K, P>(
self,
keys_of: K,
project: P,
) -> Result<Vec<Row<'static>>, EngineError>
where
K: Fn(&Row<'static>) -> Result<Vec<OrderKey>, EngineError>,
P: Fn(&Row<'static>) -> Result<Row<'static>, EngineError>,
{
let mut out: Vec<Row<'static>> = Vec::new();
self.finish_each(
keys_of,
|src, buf| {
buf.extend(project(src)?.values);
Ok(())
},
|cells| {
out.push(Row::new(cells.to_vec()));
Ok(())
},
)?;
Ok(out)
}
fn next_row<K>(
run: &mut dyn TempRun,
schema: &TableSchema,
keys_of: &K,
needed: &[bool],
) -> Result<Option<(Vec<OrderKey>, Row<'static>)>, EngineError>
where
K: Fn(&Row<'static>) -> Result<Vec<OrderKey>, EngineError>,
{
let Some(len_bytes) = read_exact(run, 4)? else {
return Ok(None);
};
let len = u32::from_le_bytes([len_bytes[0], len_bytes[1], len_bytes[2], len_bytes[3]]);
let Some(body) = read_exact(run, len as usize)? else {
return Err(EngineError::Internal(alloc::string::String::from(
"temp run ended before its row body",
)));
};
let (row, _) = spg_storage::decode_row_body_dense_pruned(
&body,
schema,
spg_storage::CURRENT_ROW_CODEC_VERSION,
needed,
)
.map_err(|e| EngineError::Internal(alloc::format!("temp run row decode: {e:?}")))?;
let keys = keys_of(&row)?;
Ok(Some((keys, row)))
}
}
#[cfg(test)]
mod tests {
use super::*;
use alloc::string::ToString;
use spg_storage::{ColumnSchema, DataType, Value};
struct MemRun {
buf: Vec<u8>,
read_at: usize,
sealed: bool,
}
impl TempRun for MemRun {
fn append(&mut self, bytes: &[u8]) -> Result<(), crate::tempstore::TempStoreError> {
assert!(!self.sealed, "append after seal");
self.buf.extend_from_slice(bytes);
Ok(())
}
fn seal(&mut self) -> Result<(), crate::tempstore::TempStoreError> {
self.sealed = true;
self.read_at = 0;
Ok(())
}
fn read(&mut self, buf: &mut [u8]) -> Result<usize, crate::tempstore::TempStoreError> {
let n = core::cmp::min(buf.len(), self.buf.len() - self.read_at);
buf[..n].copy_from_slice(&self.buf[self.read_at..self.read_at + n]);
self.read_at += n;
Ok(n)
}
fn bytes_written(&self) -> u64 {
self.buf.len() as u64
}
}
fn mem_run() -> Result<Box<dyn TempRun>, crate::tempstore::TempStoreError> {
Ok(Box::new(MemRun {
buf: Vec::new(),
read_at: 0,
sealed: false,
}))
}
fn record_cols() -> Vec<ColumnSchema> {
alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new("pad", DataType::Text, false),
]
}
fn record(id: i32) -> Row<'static> {
Row::new(alloc::vec![
Value::Int(id),
Value::text(alloc::format!("pad-{id}")),
])
}
fn keys_of(src: &Row<'static>) -> Result<Vec<OrderKey>, EngineError> {
match &src.values[0] {
Value::Int(n) => Ok(alloc::vec![OrderKey::Int(i128::from(*n))]),
other => Err(EngineError::Internal(alloc::format!("bad key: {other:?}"))),
}
}
fn ids_of(rows: &[Row<'static>], at: usize) -> Vec<i32> {
rows.iter()
.map(|r| match r.values[at] {
Value::Int(n) => n,
_ => panic!("int column at {at}"),
})
.collect()
}
fn project_pad_only(src: &Row<'static>) -> Result<Row<'static>, EngineError> {
Ok(Row::new(alloc::vec![src.values[1].clone()]))
}
fn project_identity(src: &Row<'static>) -> Result<Row<'static>, EngineError> {
Ok(src.clone())
}
#[test]
fn pruned_sort_stores_only_what_the_caller_reads() {
let descs = [false];
let cols = alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new("pad", DataType::Text, false),
];
let wide = |id: i32| {
Row::new(alloc::vec![
Value::Int(id),
Value::text(core::iter::repeat_n('x', 200).collect::<alloc::string::String>()),
])
};
let mut full = ExternalSorter::new(None, usize::MAX, cols.clone(), &descs);
let mask = [true, false];
let mut lean = ExternalSorter::new(None, usize::MAX, cols, &descs).with_pruned(&mask);
for i in 0..1000i32 {
let r = wide((i * 7919) % 1000);
let mut k = keys_of(&r).unwrap();
full.push(&mut k, &r).unwrap();
let mut k = keys_of(&r).unwrap();
lean.push(&mut k, &r).unwrap();
}
let (full_bytes, lean_bytes) = (full.arena.len(), lean.arena.len());
assert!(
full_bytes > 200 * 1000,
"control must be carrying the payload, got {full_bytes} B for 1000 rows"
);
assert!(
lean_bytes * 20 < full_bytes,
"pruned batch {lean_bytes} B is not decisively smaller than {full_bytes} B"
);
let mut seen: Vec<i32> = Vec::new();
let n = lean
.finish_each(
keys_of,
|src, buf| {
assert!(
matches!(src.values[1], Value::Null),
"a pruned column must read NULL, not {:?}",
src.values[1]
);
buf.push(src.values[0].clone());
Ok(())
},
|cells| {
match cells[0] {
Value::Int(v) => seen.push(v),
ref other => panic!("int expected, got {other:?}"),
}
Ok(())
},
)
.unwrap();
assert_eq!(n, 1000);
assert_eq!(seen, (0..1000i32).collect::<Vec<_>>(), "sorted by id");
}
#[test]
fn pruned_sort_writes_less_to_its_runs() {
let descs = [false];
let cols = alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new("pad", DataType::Text, false),
];
let wide = |id: i32| {
Row::new(alloc::vec![
Value::Int(id),
Value::text(core::iter::repeat_n('x', 200).collect::<alloc::string::String>()),
])
};
let mask = [true, false];
let run_one = |needed: &[bool]| -> (u64, u64, Vec<i32>) {
let stats = crate::tempstore::SpillStats::default();
let mut s = ExternalSorter::new(Some(mem_run), 4096, cols.clone(), &descs)
.with_stats(&stats)
.with_pruned(needed);
for i in 0..2000i32 {
let r = wide((i * 7919) % 2000);
let mut k = keys_of(&r).unwrap();
s.push(&mut k, &r).unwrap();
}
assert!(s.spilled(), "a 4 KB budget over 2000 wide rows must spill");
let mut seen: Vec<i32> = Vec::new();
s.finish_each(
keys_of,
|src, buf| {
buf.push(src.values[0].clone());
Ok(())
},
|cells| {
match cells[0] {
Value::Int(v) => seen.push(v),
ref other => panic!("int expected, got {other:?}"),
}
Ok(())
},
)
.unwrap();
(
stats.bytes.load(core::sync::atomic::Ordering::Relaxed),
stats.files.load(core::sync::atomic::Ordering::Relaxed),
seen,
)
};
let (full_bytes, full_files, full_ids) = run_one(&[]);
let (lean_bytes, lean_files, lean_ids) = run_one(&mask);
assert!(
lean_bytes * 20 < full_bytes,
"pruned run wrote {lean_bytes} B against the control's {full_bytes} B"
);
assert!(
lean_files < full_files,
"fewer bytes must mean fewer runs: {lean_files} against {full_files}"
);
assert_eq!(full_ids, lean_ids, "pruning must not change the order");
assert_eq!(lean_ids, (0..2000i32).collect::<Vec<_>>());
}
#[test]
fn finish_each_stops_when_the_consumer_stops() {
use core::cell::Cell;
let descs = [false];
let mut s = ExternalSorter::new(Some(mem_run), 4096, record_cols(), &descs);
for i in 0..5000i32 {
let r = record((i * 7919) % 5000);
let mut k = keys_of(&r).unwrap();
s.push(&mut k, &r).unwrap();
}
let projected = Cell::new(0usize);
let seen = Cell::new(0usize);
let err = s
.finish_each(
keys_of,
|src, buf| {
projected.set(projected.get() + 1);
buf.extend(project_identity(src)?.values);
Ok(())
},
|_cells| {
seen.set(seen.get() + 1);
if seen.get() == 3 {
return Err(EngineError::Internal("consumer stopped".to_string()));
}
Ok(())
},
)
.unwrap_err();
assert!(matches!(err, EngineError::Internal(ref m) if m == "consumer stopped"));
assert_eq!(seen.get(), 3, "the consumer saw exactly three rows");
assert_eq!(
projected.get(),
3,
"the merge stopped with the consumer; a collecting finish would \
have projected all 5000 first"
);
}
#[test]
fn finish_each_matches_finish_spilled_and_unspilled() {
let descs = [false];
for budget in [4096usize, 64 * 1024 * 1024] {
let mut a = ExternalSorter::new(Some(mem_run), budget, record_cols(), &descs);
let mut b = ExternalSorter::new(Some(mem_run), budget, record_cols(), &descs);
for i in 0..2000i32 {
let r = record((i * 7919) % 2000);
a.push(&mut keys_of(&r).unwrap(), &r).unwrap();
b.push(&mut keys_of(&r).unwrap(), &r).unwrap();
}
let spilled = !a.runs.is_empty();
assert_eq!(
spilled,
budget == 4096,
"budget {budget} should decide whether this spills"
);
let collected = a.finish(keys_of, project_pad_only).unwrap();
let mut streamed: Vec<Row<'static>> = Vec::new();
let n = b
.finish_each(
keys_of,
|src, buf| {
buf.extend(project_pad_only(src)?.values);
Ok(())
},
|cells| {
streamed.push(Row::new(cells.to_vec()));
Ok(())
},
)
.unwrap();
assert_eq!(n, 2000);
assert_eq!(collected, streamed, "budget {budget}");
}
}
#[test]
#[ignore]
fn r855_merge_cost_per_row_against_run_count() {
extern crate std;
use spg_storage::{ColumnSchema, DataType, Value};
use std::time::Instant;
let pad: alloc::string::String = "y".repeat(200);
let cols = alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new("pad", DataType::Text, false),
];
let descs = [false];
for rows in [50_000usize, 100_000, 200_000, 400_000] {
let mut s = ExternalSorter::new(Some(mem_run), 4 * 1024 * 1024, cols.clone(), &descs);
for i in 0..rows {
let r = Row::new(alloc::vec![
Value::Int(i32::try_from((i * 7919) % rows).unwrap()),
Value::text(pad.clone()),
]);
let mut k = keys_of(&r).unwrap();
s.push(&mut k, &r).unwrap();
}
let runs = s.runs.len() + usize::from(!s.ends.is_empty());
let t = Instant::now();
let out = s.finish(keys_of, |src| Ok(src.clone())).unwrap();
let el = t.elapsed();
std::eprintln!(
"R855 rows={rows} runs={runs} merge={el:?} per_row={:.1}ns",
el.as_nanos() as f64 / rows as f64
);
assert_eq!(out.len(), rows);
}
}
#[test]
#[ignore]
fn r853_price_the_head_scan() {
extern crate std;
use std::time::Instant;
const ROWS: usize = 400_000;
const RUNS: usize = 29;
let descs = [false];
let heads: Vec<Vec<OrderKey>> = (0..RUNS)
.map(|i| alloc::vec![OrderKey::Int(i128::try_from(i * 7919).unwrap())])
.collect();
let t = Instant::now();
let mut sink = 0usize;
for _ in 0..ROWS {
let mut best = 0usize;
for (i, k) in heads.iter().enumerate().skip(1) {
if cmp_multi_key_in(k, &heads[best], &descs, &[]) == core::cmp::Ordering::Less {
best = i;
}
}
sink += best;
}
let linear = t.elapsed();
let depth = usize::BITS as usize - RUNS.leading_zeros() as usize;
let t = Instant::now();
let mut sink2 = 0usize;
for _ in 0..ROWS {
let mut best = 0usize;
for step in 0..depth {
let i = (step + 1) % RUNS;
if cmp_multi_key_in(&heads[i], &heads[best], &descs, &[])
== core::cmp::Ordering::Less
{
best = i;
}
}
sink2 += best;
}
let treeish = t.elapsed();
std::eprintln!(
"R853 linear({RUNS} heads)={linear:?} tree-depth({depth})={treeish:?} \
cmps={} vs {} sink={sink}/{sink2}",
ROWS * (RUNS - 1),
ROWS * depth
);
}
#[test]
#[ignore]
fn r852_price_the_decode() {
extern crate std;
use spg_storage::{ColumnSchema, DataType, Value};
use std::time::Instant;
const ROWS: usize = 400_000;
let pad: alloc::string::String = "y".repeat(200);
let schema = TableSchema::new(
"spg_sort_run",
alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new("pad", DataType::Text, false),
],
);
let mut tape: Vec<u8> = Vec::new();
for i in 0..ROWS {
let row = Row::new(alloc::vec![
Value::Int(i32::try_from(i).unwrap()),
Value::text(pad.clone()),
]);
let body = spg_storage::encode_row_body_dense(&row, &schema);
tape.extend_from_slice(&u32::try_from(body.len()).unwrap().to_le_bytes());
tape.extend_from_slice(&body);
}
std::eprintln!("R852 tape={} bytes for {ROWS} rows", tape.len());
let t = Instant::now();
let mut at = 0usize;
let mut decoded = 0usize;
while at < tape.len() {
let len =
u32::from_le_bytes([tape[at], tape[at + 1], tape[at + 2], tape[at + 3]]) as usize;
at += 4;
let (row, _) = spg_storage::decode_row_body_dense(
&tape[at..at + len],
&schema,
spg_storage::CURRENT_ROW_CODEC_VERSION,
)
.unwrap();
core::hint::black_box(&row);
at += len;
decoded += 1;
}
let with_decode = t.elapsed();
let t = Instant::now();
let mut at = 0usize;
let mut seen = 0usize;
while at < tape.len() {
let len =
u32::from_le_bytes([tape[at], tape[at + 1], tape[at + 2], tape[at + 3]]) as usize;
at += 4;
core::hint::black_box(&tape[at..at + len]);
at += len;
seen += 1;
}
let walk_only = t.elapsed();
std::eprintln!(
"R852 decode+walk={with_decode:?} walk_only={walk_only:?} rows={decoded}/{seen}"
);
assert_eq!(decoded, ROWS);
assert_eq!(seen, ROWS);
}
#[test]
#[ignore]
fn r851_price_key_rederivation_in_the_merge() {
extern crate std;
use spg_storage::{ColumnSchema, DataType, Value};
use std::time::Instant;
const ROWS: usize = 400_000;
let pad: alloc::string::String = "y".repeat(200);
let cols = alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new("pad", DataType::Text, false),
];
let descs = [false];
fn free_keys(src: &Row<'static>) -> Result<Vec<OrderKey>, EngineError> {
match &src.values[0] {
Value::Int(n) => Ok(alloc::vec![OrderKey::Int(i128::from(*n))]),
other => Err(EngineError::Internal(alloc::format!("bad key: {other:?}"))),
}
}
fn computed_keys(src: &Row<'static>) -> Result<Vec<OrderKey>, EngineError> {
let Value::Int(n) = &src.values[0] else {
return Err(EngineError::Internal(alloc::string::String::from(
"bad key",
)));
};
let Value::Text(t) = &src.values[1] else {
return Err(EngineError::Internal(alloc::string::String::from(
"bad pad",
)));
};
let folded = i128::from(*n) + i128::try_from(t.len()).unwrap_or(0);
Ok(alloc::vec![OrderKey::Int(folded)])
}
for (label, keyfn, narrow) in [
(
"free",
free_keys as fn(&Row<'static>) -> Result<Vec<OrderKey>, EngineError>,
false,
),
("computed", computed_keys, false),
("narrow-projection", free_keys, true),
] {
let mut s = ExternalSorter::new(Some(mem_run), 4 * 1024 * 1024, cols.clone(), &descs);
for i in 0..ROWS {
let r = Row::new(alloc::vec![
Value::Int(i32::try_from((i * 7919) % ROWS).unwrap()),
Value::text(pad.clone()),
]);
let mut k = keyfn(&r).unwrap();
s.push(&mut k, &r).unwrap();
}
let t = Instant::now();
let out = if narrow {
s.finish(keyfn, |src| {
Ok(Row::new(alloc::vec![src.values[0].clone()]))
})
} else {
s.finish(keyfn, |src| Ok(src.clone()))
}
.unwrap();
std::eprintln!("R851 {label} merge={:?} rows={}", t.elapsed(), out.len());
assert_eq!(out.len(), ROWS);
}
}
#[test]
#[ignore]
fn r844_merge_cost_against_run_count() {
extern crate std;
use spg_storage::{ColumnSchema, DataType, Value};
use std::time::Instant;
const ROWS: usize = 400_000;
let pad: alloc::string::String = "y".repeat(200);
let cols = alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new("pad", DataType::Text, false),
];
let descs = [false];
for budget_mb in [4usize, 16, 64] {
let mut s =
ExternalSorter::new(Some(mem_run), budget_mb * 1024 * 1024, cols.clone(), &descs);
for i in 0..ROWS {
let r = Row::new(alloc::vec![
Value::Int(i32::try_from((i * 7919) % ROWS).unwrap()),
Value::text(pad.clone()),
]);
let mut k = keys_of(&r).unwrap();
s.push(&mut k, &r).unwrap();
}
let runs = s.runs.len() + usize::from(!s.ends.is_empty());
let t = Instant::now();
let out = s.finish(keys_of, |src| Ok(src.clone())).unwrap();
std::eprintln!(
"R844 budget={budget_mb}MB runs={runs} merge={:?} rows={}",
t.elapsed(),
out.len()
);
assert_eq!(out.len(), ROWS);
}
}
#[test]
#[ignore]
fn r843_price_the_wirings_extra_work() {
extern crate std;
use spg_storage::{ColumnSchema, DataType, Value};
use std::time::Instant;
const ROWS: usize = 400_000;
let pad: alloc::string::String = "y".repeat(200);
let rows: Vec<Row<'static>> = (0..ROWS)
.map(|i| {
Row::new(alloc::vec![
Value::Int(i32::try_from(i).unwrap()),
Value::text(pad.clone()),
])
})
.collect();
let t0 = Instant::now();
let mut sink: Vec<Row<'static>> = Vec::with_capacity(ROWS);
for r in &rows {
sink.push(r.clone());
}
let clone_cost = t0.elapsed();
std::hint::black_box(&sink);
let t1 = Instant::now();
let out: Vec<Row<'static>> = sink.iter().map(Clone::clone).collect();
let second_pass = t1.elapsed();
std::hint::black_box(&out);
let cols = alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new("pad", DataType::Text, false),
];
std::hint::black_box(&cols);
std::eprintln!(
"R843 rows={ROWS} source_row_clone={clone_cost:?} extra_pass={second_pass:?}"
);
}
#[test]
#[ignore]
fn r839_profile_spill_phases() {
extern crate std;
use spg_storage::{ColumnSchema, DataType, Value};
use std::time::Instant;
const ROWS: usize = 400_000;
const WORK_MEM: usize = 4 * 1024 * 1024;
let cols = alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new("pad", DataType::Text, false),
];
let pad: alloc::string::String = "y".repeat(200);
let descs = [false];
let t0 = Instant::now();
let rows: Vec<Row<'static>> = (0..ROWS)
.map(|i| {
Row::new(alloc::vec![
Value::Int(i32::try_from((i * 7919) % ROWS).unwrap()),
Value::text(pad.clone()),
])
})
.collect();
let build = t0.elapsed();
let mut s = ExternalSorter::new(Some(mem_run), WORK_MEM, cols, &descs);
let t1 = Instant::now();
for r in rows {
let mut k = keys_of(&r).unwrap();
s.push(&mut k, &r).unwrap();
}
let push_phase = t1.elapsed();
let t2 = Instant::now();
let out = s.finish(keys_of, |src| Ok(src.clone())).unwrap();
let merge_phase = t2.elapsed();
std::eprintln!(
"R839 rows={ROWS} build={build:?} push_spill={push_phase:?} merge_project={merge_phase:?} out={}",
out.len()
);
assert_eq!(out.len(), ROWS);
}
#[test]
fn a_key_that_was_never_projected_still_orders_the_merge() {
let descs = [false];
let mut s = ExternalSorter::new(Some(mem_run), 1, record_cols(), &descs);
for id in [7, 3, 9, 1, 8, 2, 6, 4, 5, 0] {
let r = record(id);
let mut k = keys_of(&r).unwrap();
s.push(&mut k, &r).unwrap();
}
assert!(s.spilled(), "a 1-byte budget must have produced runs");
let out = s.finish(keys_of, project_pad_only).unwrap();
let pads: Vec<alloc::string::String> = out
.iter()
.map(|r| match &r.values[0] {
Value::Text(t) => t.to_string(),
other => panic!("pad column, got {other:?}"),
})
.collect();
assert_eq!(
pads,
alloc::vec![
"pad-0", "pad-1", "pad-2", "pad-3", "pad-4", "pad-5", "pad-6", "pad-7", "pad-8",
"pad-9"
],
"ordered by the unprojected key"
);
assert_eq!(
out[0].values.len(),
1,
"the key column must not leak into the result"
);
}
#[test]
fn a_spilled_sort_returns_exactly_what_an_in_memory_sort_would() {
let descs = [false];
let mut s = ExternalSorter::new(Some(mem_run), 1, record_cols(), &descs);
for id in [7, 3, 9, 1, 8, 2, 6, 4, 5, 0] {
let r = record(id);
let mut k = keys_of(&r).unwrap();
s.push(&mut k, &r).unwrap();
}
assert!(s.spilled());
let out = s.finish(keys_of, project_identity).unwrap();
assert_eq!(ids_of(&out, 0), alloc::vec![0, 1, 2, 3, 4, 5, 6, 7, 8, 9]);
assert_eq!(
out[0].values[1],
Value::text("pad-0".to_string()),
"the whole row comes back, not only its sort key"
);
}
#[test]
fn descending_order_survives_the_merge() {
let descs = [true];
let mut s = ExternalSorter::new(Some(mem_run), 1, record_cols(), &descs);
for id in [3, 1, 2] {
let r = record(id);
let mut k = keys_of(&r).unwrap();
s.push(&mut k, &r).unwrap();
}
assert!(s.spilled());
assert_eq!(
ids_of(&s.finish(keys_of, project_identity).unwrap(), 0),
alloc::vec![3, 2, 1],
"DESC must survive the k-way merge"
);
}
#[test]
fn without_a_factory_the_sort_stays_in_memory_and_still_sorts() {
let descs = [false];
let mut s = ExternalSorter::new(None, 1, record_cols(), &descs);
for id in [2, 0, 1] {
let r = record(id);
let mut k = keys_of(&r).unwrap();
s.push(&mut k, &r).unwrap();
}
assert!(!s.spilled(), "nothing to spill to means nothing spilled");
let out = s.finish(keys_of, project_pad_only).unwrap();
let pads: Vec<alloc::string::String> = out
.iter()
.map(|r| match &r.values[0] {
Value::Text(t) => t.to_string(),
other => panic!("pad column, got {other:?}"),
})
.collect();
assert_eq!(pads, alloc::vec!["pad-0", "pad-1", "pad-2"]);
assert_eq!(
out[0].values.len(),
1,
"same output shape as the spilled path"
);
}
fn assert_order_matches(name: &str, keyrows: &[Vec<OrderKey>], descs: &[bool]) {
let mut s = ExternalSorter::new(None, usize::MAX, record_cols(), descs);
for (i, k) in keyrows.iter().enumerate() {
let r = record(i as i32);
s.push(&mut k.clone(), &r).unwrap();
}
let mut want: Vec<u32> = (0..keyrows.len() as u32).collect();
want.sort_by(|&a, &b| {
cmp_multi_key_in(&keyrows[a as usize], &keyrows[b as usize], descs, &[])
});
assert_eq!(s.sorted_order(), want, "{name}");
}
#[test]
fn inline_int_key_orders_the_batch_like_the_general_comparator() {
let int = |n: i128| alloc::vec![OrderKey::Int(n)];
let plain: Vec<Vec<OrderKey>> = alloc::vec![
int(5),
int(-3),
int(5),
int(0),
int(i64::MAX as i128),
int(-1),
int(5),
];
let with_nulls: Vec<Vec<OrderKey>> = alloc::vec![
int(7),
alloc::vec![OrderKey::NullBig],
int(-2),
alloc::vec![OrderKey::NullSmall],
int(7),
alloc::vec![OrderKey::NullBig],
];
let at_the_ends: Vec<Vec<OrderKey>> =
alloc::vec![int(3), int(i128::MAX), int(-4), int(i128::MIN), int(3)];
let texts: Vec<Vec<OrderKey>> = alloc::vec![
alloc::vec![OrderKey::Text("pear".into())],
alloc::vec![OrderKey::Text("apple".into())],
alloc::vec![OrderKey::Num(1.5)],
];
let two_keys: Vec<Vec<OrderKey>> = alloc::vec![
alloc::vec![OrderKey::Int(1), OrderKey::Int(9)],
alloc::vec![OrderKey::Int(1), OrderKey::Int(2)],
alloc::vec![OrderKey::Int(0), OrderKey::Int(5)],
];
for (name, rows) in [
("plain ints", &plain),
("null sentinels", &with_nulls),
("keys at the ends of the range", &at_the_ends),
("non-integer keys", &texts),
] {
assert_order_matches(name, rows, &[false]);
assert_order_matches(name, rows, &[true]);
}
assert_order_matches("two keys", &two_keys, &[false, false]);
assert_order_matches("two keys, second descending", &two_keys, &[false, true]);
}
#[test]
fn inline_int_key_holds_when_the_sort_spills() {
let descs = [false];
for budget in [4096usize, 64 * 1024 * 1024] {
let mut s = ExternalSorter::new(Some(mem_run), budget, record_cols(), &descs);
for i in 0..2000i32 {
let r = record((i * 7919) % 2000);
s.push(&mut keys_of(&r).unwrap(), &r).unwrap();
}
assert_eq!(s.spilled(), budget == 4096, "budget {budget}");
let out = s.finish(keys_of, project_identity).unwrap();
assert_eq!(
ids_of(&out, 0),
(0..2000i32).collect::<Vec<i32>>(),
"budget {budget}"
);
}
}
}