use super::*;
pub struct LimitedResult {
inner: Box<dyn QueryResult>,
limit: Option<usize>,
offset: usize,
returned_count: usize,
offset_applied: bool,
columns: Vec<String>,
}
impl LimitedResult {
pub fn new(inner: Box<dyn QueryResult>, limit: Option<usize>, offset: usize) -> Self {
let columns = inner.columns().to_vec();
Self {
inner,
limit,
offset,
returned_count: 0,
offset_applied: false,
columns,
}
}
pub fn with_limit(inner: Box<dyn QueryResult>, limit: usize) -> Self {
Self::new(inner, Some(limit), 0)
}
pub fn with_offset(inner: Box<dyn QueryResult>, offset: usize) -> Self {
Self::new(inner, None, offset)
}
}
impl QueryResult for LimitedResult {
fn columns(&self) -> &[String] {
&self.columns
}
fn columns_arc(&self) -> Option<CompactArc<Vec<String>>> {
self.inner.columns_arc()
}
fn next(&mut self) -> bool {
if !self.offset_applied {
for _ in 0..self.offset {
if !self.inner.next() {
self.offset_applied = true;
return false;
}
}
self.offset_applied = true;
}
if let Some(limit) = self.limit {
if self.returned_count >= limit {
return false;
}
}
if self.inner.next() {
self.returned_count += 1;
true
} else {
false
}
}
fn scan(&self, dest: &mut [Value]) -> Result<()> {
self.inner.scan(dest)
}
fn row(&self) -> &Row {
self.inner.row()
}
fn take_row(&mut self) -> Row {
self.inner.take_row()
}
fn take_deferred_row(&mut self) -> DeferredRow {
self.inner.take_deferred_row()
}
fn preserves_deferred_rows(&self) -> bool {
self.inner.preserves_deferred_rows()
}
fn ascending_nulls_last_ordering(&self) -> Option<Vec<usize>> {
self.inner.ascending_nulls_last_ordering()
}
fn close(&mut self) -> Result<()> {
self.inner.close()
}
fn rows_affected(&self) -> i64 {
self.inner.rows_affected()
}
fn last_insert_id(&self) -> i64 {
self.inner.last_insert_id()
}
fn last_error(&mut self) -> Option<radixdb_core::Error> {
self.inner.last_error()
}
fn estimated_count(&self) -> Option<usize> {
self.inner.estimated_count().map(|rows| {
let after_offset = if self.offset_applied {
rows
} else {
rows.saturating_sub(self.offset)
};
self.limit.map_or(after_offset, |limit| {
after_offset.min(limit.saturating_sub(self.returned_count))
})
})
}
fn with_aliases(self: Box<Self>, aliases: FxHashMap<String, String>) -> Box<dyn QueryResult> {
Box::new(AliasedResult::new(self, aliases))
}
}
pub struct OrderedResult {
inner: Box<dyn QueryResult>,
}
pub(super) const ORDERED_RUN_MAX_ROWS: usize = 131_072;
pub(super) const ORDERED_RUN_MAX_BYTES: usize = 64 * 1024 * 1024;
const ORDERED_SPILL_MAX_VALUE_BYTES: usize = 64 * 1024 * 1024;
static ORDERED_SPILL_ID: AtomicU64 = AtomicU64::new(1);
type OrderedRowComparator = dyn Fn(&Row, &Row) -> std::cmp::Ordering + Send;
enum BoundedOrderedRows {
Memory {
rows: RowVec,
input_rows: usize,
peak_rows: usize,
peak_bytes: usize,
},
External {
paths: Vec<PathBuf>,
input_rows: usize,
peak_rows: usize,
peak_bytes: usize,
},
}
struct OrderedSpillRun {
path: PathBuf,
reader: BufReader<File>,
remaining: u64,
head: Option<Row>,
}
impl OrderedSpillRun {
fn open(path: PathBuf) -> Result<Self> {
let opened = (|| {
let file = File::open(&path).map_err(|error| {
Error::internal(format!(
"failed to open ORDER BY spill run {}: {error}",
path.display()
))
})?;
let mut reader = BufReader::new(file);
let remaining = read_u64(&mut reader, "ORDER BY spill row count")?;
let mut run = Self {
path: path.clone(),
reader,
remaining,
head: None,
};
run.advance()?;
Ok(run)
})();
if opened.is_err() {
let _ = std::fs::remove_file(path);
}
opened
}
fn advance(&mut self) -> Result<()> {
self.head = if self.remaining == 0 {
None
} else {
self.remaining -= 1;
Some(read_spill_row(&mut self.reader)?)
};
Ok(())
}
}
impl Drop for OrderedSpillRun {
fn drop(&mut self) {
let _ = std::fs::remove_file(&self.path);
}
}
struct ExternalOrderedResult {
columns: CompactArc<Vec<String>>,
runs: Vec<OrderedSpillRun>,
compare: Box<OrderedRowComparator>,
current: Option<Row>,
remaining: usize,
closed: bool,
last_error: Option<Error>,
}
impl ExternalOrderedResult {
fn new<F>(
columns: Vec<String>,
paths: Vec<PathBuf>,
input_rows: usize,
compare: F,
) -> Result<Self>
where
F: Fn(&Row, &Row) -> std::cmp::Ordering + Send + 'static,
{
let mut pending = paths.into_iter();
let mut runs = Vec::new();
while let Some(path) = pending.next() {
match OrderedSpillRun::open(path) {
Ok(run) => runs.push(run),
Err(error) => {
for path in pending {
let _ = std::fs::remove_file(path);
}
return Err(error);
}
}
}
Ok(Self {
columns: CompactArc::new(columns),
runs,
compare: Box::new(compare),
current: None,
remaining: input_rows,
closed: false,
last_error: None,
})
}
}
impl QueryResult for ExternalOrderedResult {
fn columns(&self) -> &[String] {
&self.columns
}
fn columns_arc(&self) -> Option<CompactArc<Vec<String>>> {
Some(CompactArc::clone(&self.columns))
}
fn next(&mut self) -> bool {
if self.closed || self.last_error.is_some() {
return false;
}
let mut best: Option<usize> = None;
for (index, run) in self.runs.iter().enumerate() {
let Some(candidate) = run.head.as_ref() else {
continue;
};
if best.is_none_or(|best_index| {
let best_row = self.runs[best_index]
.head
.as_ref()
.expect("selected ORDER BY run has a head row");
(self.compare)(candidate, best_row).is_lt()
}) {
best = Some(index);
}
}
let Some(best) = best else {
self.current = None;
return false;
};
self.current = self.runs[best].head.take();
if let Err(error) = self.runs[best].advance() {
self.last_error = Some(error);
self.current = None;
return false;
}
self.remaining = self.remaining.saturating_sub(1);
true
}
fn scan(&self, dest: &mut [Value]) -> Result<()> {
let row = self.row();
if dest.len() != row.len() {
return Err(Error::internal(format!(
"scan destination has {} values for ORDER BY row with {} columns",
dest.len(),
row.len()
)));
}
dest.clone_from_slice(row.as_slice());
Ok(())
}
fn row(&self) -> &Row {
self.current
.as_ref()
.expect("row() called without successful external ORDER BY next()")
}
fn take_row(&mut self) -> Row {
self.current
.take()
.expect("take_row() called without successful external ORDER BY next()")
}
fn close(&mut self) -> Result<()> {
self.closed = true;
self.current = None;
self.runs.clear();
Ok(())
}
fn rows_affected(&self) -> i64 {
0
}
fn last_insert_id(&self) -> i64 {
0
}
fn last_error(&mut self) -> Option<Error> {
self.last_error.take()
}
fn estimated_count(&self) -> Option<usize> {
Some(self.remaining)
}
fn with_aliases(self: Box<Self>, aliases: FxHashMap<String, String>) -> Box<dyn QueryResult> {
Box::new(AliasedResult::new(self, aliases))
}
}
#[derive(Clone, Copy)]
pub struct RadixOrderSpec {
pub col_idx: usize,
pub ascending: bool,
pub nulls_first: Option<bool>,
}
fn ordered_spill_path() -> Result<PathBuf> {
let base = std::env::temp_dir();
for _ in 0..32 {
let ordinal = ORDERED_SPILL_ID.fetch_add(1, AtomicOrdering::Relaxed);
let path = base.join(format!(
"radixdb-order-{}-{ordinal}.run",
std::process::id()
));
match OpenOptions::new().write(true).create_new(true).open(&path) {
Ok(file) => {
drop(file);
return Ok(path);
}
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => continue,
Err(error) => {
return Err(Error::internal(format!(
"failed to create ORDER BY spill path {}: {error}",
path.display()
)))
}
}
}
Err(Error::internal(
"failed to allocate a unique ORDER BY spill path",
))
}
fn write_u32(writer: &mut impl Write, value: usize, what: &str) -> Result<()> {
let value =
u32::try_from(value).map_err(|_| Error::invalid_argument(format!("{what} exceeds u32")))?;
writer
.write_all(&value.to_le_bytes())
.map_err(|error| Error::internal(format!("failed to write {what}: {error}")))
}
fn write_spill_row(writer: &mut impl Write, row: &Row) -> Result<()> {
write_u32(writer, row.len(), "ORDER BY spill column count")?;
let mut encoded = Vec::new();
for value in row.iter() {
encoded.clear();
radixdb_storage::mvcc::persistence::serialize_value_into(&mut encoded, value)?;
write_u32(writer, encoded.len(), "ORDER BY spill value length")?;
writer.write_all(&encoded).map_err(|error| {
Error::internal(format!("failed to write ORDER BY spill value: {error}"))
})?;
}
Ok(())
}
fn read_u32(reader: &mut impl Read, what: &str) -> Result<u32> {
let mut bytes = [0u8; 4];
reader
.read_exact(&mut bytes)
.map_err(|error| Error::internal(format!("failed to read {what}: {error}")))?;
Ok(u32::from_le_bytes(bytes))
}
fn read_u64(reader: &mut impl Read, what: &str) -> Result<u64> {
let mut bytes = [0u8; 8];
reader
.read_exact(&mut bytes)
.map_err(|error| Error::internal(format!("failed to read {what}: {error}")))?;
Ok(u64::from_le_bytes(bytes))
}
fn read_spill_row(reader: &mut impl Read) -> Result<Row> {
let columns = read_u32(reader, "ORDER BY spill column count")? as usize;
if columns > RetainedRowsBudget::DEFAULT_MAX_ROWS {
return Err(Error::internal(format!(
"ORDER BY spill row declares {columns} columns"
)));
}
let mut row = Row::with_capacity(columns);
for _ in 0..columns {
let value_len = read_u32(reader, "ORDER BY spill value length")? as usize;
if value_len > ORDERED_SPILL_MAX_VALUE_BYTES {
return Err(Error::internal(format!(
"ORDER BY spill value length {value_len} exceeds bounded decoder limit"
)));
}
let mut encoded = vec![0u8; value_len];
reader.read_exact(&mut encoded).map_err(|error| {
Error::internal(format!("failed to read ORDER BY spill value: {error}"))
})?;
row.push(radixdb_storage::mvcc::persistence::deserialize_value(
&encoded,
)?);
}
Ok(row)
}
fn write_ordered_run<F>(rows: &mut RowVec, compare: &F) -> Result<PathBuf>
where
F: Fn(&Row, &Row) -> std::cmp::Ordering,
{
rows.sort_unstable_by(|(_, left), (_, right)| compare(left, right));
let path = ordered_spill_path()?;
let result = (|| {
let file = OpenOptions::new()
.write(true)
.truncate(true)
.open(&path)
.map_err(|error| {
Error::internal(format!(
"failed to open ORDER BY spill run {}: {error}",
path.display()
))
})?;
let mut writer = BufWriter::new(file);
writer
.write_all(&(rows.len() as u64).to_le_bytes())
.map_err(|error| Error::internal(format!("failed to write spill header: {error}")))?;
for (_, row) in rows.iter() {
write_spill_row(&mut writer, row)?;
}
writer
.flush()
.map_err(|error| Error::internal(format!("failed to flush ORDER BY spill: {error}")))
})();
if result.is_err() {
let _ = std::fs::remove_file(&path);
}
result.map(|()| path)
}
fn remove_ordered_runs(paths: impl IntoIterator<Item = PathBuf>) {
for path in paths {
let _ = std::fs::remove_file(path);
}
}
fn collect_bounded_ordered_rows<F>(
mut inner: Box<dyn QueryResult>,
compare: &F,
) -> Result<(Vec<String>, BoundedOrderedRows)>
where
F: Fn(&Row, &Row) -> std::cmp::Ordering,
{
let columns = inner.columns().to_vec();
let mut rows = RowVec::with_capacity(ORDERED_RUN_MAX_ROWS.min(1024));
let mut budget = RetainedRowsBudget::with_limits(
"bounded ORDER BY run",
ORDERED_RUN_MAX_ROWS,
ORDERED_RUN_MAX_BYTES,
);
let mut paths = Vec::new();
let mut input_rows = 0usize;
let mut peak_rows = 0usize;
let mut peak_bytes = 0usize;
while inner.next() {
let row = inner.take_row();
if let Err(error) = budget.admit(&row) {
if rows.is_empty() {
return Err(error);
}
peak_rows = peak_rows.max(budget.peak_rows());
peak_bytes = peak_bytes.max(budget.peak_bytes());
let path = match write_ordered_run(&mut rows, compare) {
Ok(path) => path,
Err(error) => {
remove_ordered_runs(paths);
return Err(error);
}
};
paths.push(path);
rows.clear();
budget = RetainedRowsBudget::with_limits(
"bounded ORDER BY run",
ORDERED_RUN_MAX_ROWS,
ORDERED_RUN_MAX_BYTES,
);
if let Err(error) = budget.admit(&row) {
remove_ordered_runs(paths);
return Err(error);
}
}
rows.push((input_rows as i64, row));
input_rows = input_rows.saturating_add(1);
}
if let Some(error) = inner.last_error() {
remove_ordered_runs(paths);
return Err(error);
}
peak_rows = peak_rows.max(budget.peak_rows());
peak_bytes = peak_bytes.max(budget.peak_bytes());
if paths.is_empty() {
return Ok((
columns,
BoundedOrderedRows::Memory {
rows,
input_rows,
peak_rows,
peak_bytes,
},
));
}
if !rows.is_empty() {
match write_ordered_run(&mut rows, compare) {
Ok(path) => paths.push(path),
Err(error) => {
remove_ordered_runs(paths);
return Err(error);
}
}
}
Ok((
columns,
BoundedOrderedRows::External {
paths,
input_rows,
peak_rows,
peak_bytes,
},
))
}
impl OrderedResult {
pub fn new<F>(inner: Box<dyn QueryResult>, compare: F) -> Result<Self>
where
F: Fn(&Row, &Row) -> std::cmp::Ordering + Send + 'static,
{
let collect_started = radixdb_core::time_compat::Instant::now();
let (columns, bounded) = collect_bounded_ordered_rows(inner, &compare)?;
let collect_elapsed = collect_started.elapsed();
let inner: Box<dyn QueryResult> = match bounded {
BoundedOrderedRows::Memory {
mut rows,
input_rows,
peak_rows,
peak_bytes,
} => {
let finalize_started = radixdb_core::time_compat::Instant::now();
rows.sort_unstable_by(|(_, left), (_, right)| compare(left, right));
let finalize_elapsed = finalize_started.elapsed();
radixdb_storage::instrumentation::record_join_ordered_sort(
input_rows as u64,
0,
peak_rows as u64,
peak_bytes as u64,
collect_elapsed,
finalize_elapsed,
);
Box::new(ExecutorResult::new(columns, rows))
}
BoundedOrderedRows::External {
paths,
input_rows,
peak_rows,
peak_bytes,
} => {
let runs = paths.len();
let finalize_started = radixdb_core::time_compat::Instant::now();
let external = ExternalOrderedResult::new(columns, paths, input_rows, compare)?;
let finalize_elapsed = finalize_started.elapsed();
radixdb_storage::instrumentation::record_join_ordered_sort(
input_rows as u64,
runs as u64,
peak_rows as u64,
peak_bytes as u64,
collect_elapsed,
finalize_elapsed,
);
Box::new(external)
}
};
Ok(Self { inner })
}
pub fn new_radix<F>(
inner: Box<dyn QueryResult>,
order_specs: &[RadixOrderSpec],
fallback_compare: F,
) -> Result<Self>
where
F: Fn(&Row, &Row) -> std::cmp::Ordering + Send + 'static,
{
let collect_started = radixdb_core::time_compat::Instant::now();
let (columns, bounded) = collect_bounded_ordered_rows(inner, &fallback_compare)?;
let collect_elapsed = collect_started.elapsed();
let BoundedOrderedRows::Memory {
mut rows,
input_rows,
peak_rows,
peak_bytes,
} = bounded
else {
let BoundedOrderedRows::External {
paths,
input_rows,
peak_rows,
peak_bytes,
} = bounded
else {
unreachable!()
};
let runs = paths.len();
let finalize_started = radixdb_core::time_compat::Instant::now();
let external =
ExternalOrderedResult::new(columns, paths, input_rows, fallback_compare)?;
let finalize_elapsed = finalize_started.elapsed();
radixdb_storage::instrumentation::record_join_ordered_sort(
input_rows as u64,
runs as u64,
peak_rows as u64,
peak_bytes as u64,
collect_elapsed,
finalize_elapsed,
);
return Ok(Self {
inner: Box::new(external),
});
};
let has_explicit_nulls_ordering = order_specs.iter().any(|s| s.nulls_first.is_some());
if !has_explicit_nulls_ordering {
if order_specs.len() == 1 {
let spec = &order_specs[0];
let finalize_started = radixdb_core::time_compat::Instant::now();
if Self::try_radix_sort_single_int(&mut rows, spec.col_idx, spec.ascending) {
let finalize_elapsed = finalize_started.elapsed();
radixdb_storage::instrumentation::record_join_ordered_sort(
input_rows as u64,
0,
peak_rows as u64,
peak_bytes as u64,
collect_elapsed,
finalize_elapsed,
);
return Ok(Self {
inner: Box::new(ExecutorResult::new(columns, rows)),
});
}
}
let finalize_started = radixdb_core::time_compat::Instant::now();
if order_specs.len() <= 4 && Self::try_radix_sort_multi_int(&mut rows, order_specs) {
let finalize_elapsed = finalize_started.elapsed();
radixdb_storage::instrumentation::record_join_ordered_sort(
input_rows as u64,
0,
peak_rows as u64,
peak_bytes as u64,
collect_elapsed,
finalize_elapsed,
);
return Ok(Self {
inner: Box::new(ExecutorResult::new(columns, rows)),
});
}
}
let finalize_started = radixdb_core::time_compat::Instant::now();
rows.sort_unstable_by(|(_, a), (_, b)| fallback_compare(a, b));
let finalize_elapsed = finalize_started.elapsed();
radixdb_storage::instrumentation::record_join_ordered_sort(
input_rows as u64,
0,
peak_rows as u64,
peak_bytes as u64,
collect_elapsed,
finalize_elapsed,
);
Ok(Self {
inner: Box::new(ExecutorResult::new(columns, rows)),
})
}
fn try_radix_sort_single_int(rows: &mut RowVec, col_idx: usize, ascending: bool) -> bool {
if Self::try_radix_sort_single_uuid(rows, col_idx, ascending) {
return true;
}
for (_, row) in rows.iter() {
match row.get(col_idx) {
Some(Value::Integer(value)) if *value != i64::MIN => continue,
_ => return false, }
}
if ascending {
radsort::sort_by_key(rows, |(_, row)| match row.get(col_idx) {
Some(Value::Integer(i)) => *i,
_ => unreachable!("radix admission requires non-null integers"),
});
} else {
radsort::sort_by_key(rows, |(_, row)| {
match row.get(col_idx) {
Some(Value::Integer(i)) => {
i.wrapping_neg().wrapping_sub(1)
}
_ => unreachable!("radix admission requires non-null integers"),
}
});
}
true
}
pub(super) fn try_radix_sort_single_uuid(
rows: &mut RowVec,
col_idx: usize,
ascending: bool,
) -> bool {
if rows
.iter()
.any(|(_, row)| row.get(col_idx).and_then(Value::as_uuid_bytes).is_none())
{
return false;
}
let word = |row: &Row, range: std::ops::Range<usize>| {
let bytes = row
.get(col_idx)
.and_then(Value::as_uuid_bytes)
.expect("UUID radix admission validates every key");
let word = u64::from_be_bytes(
bytes[range]
.try_into()
.expect("UUID radix word is exactly eight bytes"),
);
if ascending {
word
} else {
!word
}
};
radsort::sort_by_key(rows, |(_, row)| word(row, 8..16));
radsort::sort_by_key(rows, |(_, row)| word(row, 0..8));
true
}
fn try_radix_sort_multi_int(rows: &mut RowVec, order_specs: &[RadixOrderSpec]) -> bool {
for (_, row) in rows.iter() {
for spec in order_specs {
match row.get(spec.col_idx) {
Some(Value::Integer(value)) if *value != i64::MIN => continue,
_ => return false,
}
}
}
for spec in order_specs.iter().rev() {
if spec.ascending {
radsort::sort_by_key(rows, |(_, row)| match row.get(spec.col_idx) {
Some(Value::Integer(i)) => *i,
_ => unreachable!("radix admission requires non-null integers"),
});
} else {
radsort::sort_by_key(rows, |(_, row)| match row.get(spec.col_idx) {
Some(Value::Integer(i)) => i.wrapping_neg().wrapping_sub(1),
_ => unreachable!("radix admission requires non-null integers"),
});
}
}
true
}
}
impl QueryResult for OrderedResult {
fn columns(&self) -> &[String] {
self.inner.columns()
}
fn columns_arc(&self) -> Option<CompactArc<Vec<String>>> {
self.inner.columns_arc()
}
fn next(&mut self) -> bool {
self.inner.next()
}
fn scan(&self, dest: &mut [Value]) -> Result<()> {
self.inner.scan(dest)
}
fn row(&self) -> &Row {
self.inner.row()
}
fn take_row(&mut self) -> Row {
self.inner.take_row()
}
fn close(&mut self) -> Result<()> {
self.inner.close()
}
fn rows_affected(&self) -> i64 {
0
}
fn last_insert_id(&self) -> i64 {
self.inner.last_insert_id()
}
fn last_error(&mut self) -> Option<Error> {
self.inner.last_error()
}
fn estimated_count(&self) -> Option<usize> {
self.inner.estimated_count()
}
fn with_aliases(self: Box<Self>, aliases: FxHashMap<String, String>) -> Box<dyn QueryResult> {
Box::new(AliasedResult::new(self, aliases))
}
}
pub struct TopNResult {
inner: ExecutorResult,
_budget: RetainedRowsBudget,
}
impl TopNResult {
pub fn new<F>(
inner: Box<dyn QueryResult>,
compare: F,
limit: usize,
offset: usize,
) -> Result<Self>
where
F: Fn(&Row, &Row) -> std::cmp::Ordering + Clone,
{
Self::build(
inner,
compare,
limit,
offset,
RetainedRowsBudget::new("TOP-N"),
)
}
pub fn new_with_context<F>(
inner: Box<dyn QueryResult>,
compare: F,
limit: usize,
offset: usize,
ctx: &crate::context::ExecutionContext,
) -> Result<Self>
where
F: Fn(&Row, &Row) -> std::cmp::Ordering + Clone,
{
Self::build(
inner,
compare,
limit,
offset,
RetainedRowsBudget::with_request_memory("TOP-N", ctx)?,
)
}
fn build<F>(
mut inner: Box<dyn QueryResult>,
compare: F,
limit: usize,
offset: usize,
mut budget: RetainedRowsBudget,
) -> Result<Self>
where
F: Fn(&Row, &Row) -> std::cmp::Ordering + Clone,
{
use std::collections::BinaryHeap;
let columns = inner.columns().to_vec();
let heap_capacity = limit.saturating_add(offset);
budget.ensure_capacity(heap_capacity)?;
if heap_capacity == 0 {
return Ok(Self {
inner: ExecutorResult::new(columns, RowVec::new()),
_budget: budget,
});
}
let compare = std::sync::Arc::new(compare);
struct HeapRow<F: Fn(&Row, &Row) -> std::cmp::Ordering> {
row: Row,
compare: std::sync::Arc<F>,
}
impl<F: Fn(&Row, &Row) -> std::cmp::Ordering> PartialEq for HeapRow<F> {
fn eq(&self, other: &Self) -> bool {
(self.compare)(&self.row, &other.row) == std::cmp::Ordering::Equal
}
}
impl<F: Fn(&Row, &Row) -> std::cmp::Ordering> Eq for HeapRow<F> {}
impl<F: Fn(&Row, &Row) -> std::cmp::Ordering> PartialOrd for HeapRow<F> {
fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
Some(self.cmp(other))
}
}
impl<F: Fn(&Row, &Row) -> std::cmp::Ordering> Ord for HeapRow<F> {
fn cmp(&self, other: &Self) -> std::cmp::Ordering {
(self.compare)(&self.row, &other.row)
}
}
let mut heap: BinaryHeap<HeapRow<F>> = BinaryHeap::with_capacity(heap_capacity + 1);
let mut input_rows = 0_u64;
while inner.next() {
input_rows = input_rows.saturating_add(1);
let row = inner.take_row();
if heap.len() < heap_capacity {
budget.admit(&row)?;
heap.push(HeapRow {
row,
compare: std::sync::Arc::clone(&compare),
});
} else if let Some(worst) = heap.peek() {
if compare(&row, &worst.row) == std::cmp::Ordering::Less {
if let Some(removed) = heap.pop() {
budget.release(&removed.row);
}
budget.admit(&row)?;
heap.push(HeapRow {
row,
compare: std::sync::Arc::clone(&compare),
});
}
}
}
if let Some(err) = inner.last_error() {
return Err(err);
}
let mut rows: Vec<Row> = heap.into_iter().map(|hr| hr.row).collect();
rows.sort_unstable_by(|a, b| compare(a, b));
if offset > 0 && offset < rows.len() {
for row in rows.drain(..offset) {
budget.release(&row);
}
} else if offset >= rows.len() {
for row in &rows {
budget.release(row);
}
rows.clear();
}
let result_rows: RowVec = rows
.into_iter()
.enumerate()
.map(|(i, row)| (i as i64, row))
.collect();
radixdb_storage::instrumentation::record_join_top_n(
input_rows,
budget.peak_rows() as u64,
budget.peak_bytes() as u64,
result_rows.len() as u64,
);
Ok(Self {
inner: ExecutorResult::new(columns, result_rows),
_budget: budget,
})
}
pub fn from_rows_with_budget(
columns: Vec<String>,
rows: RowVec,
budget: RetainedRowsBudget,
) -> Self {
Self {
inner: ExecutorResult::new(columns, rows),
_budget: budget,
}
}
}
impl QueryResult for TopNResult {
fn columns(&self) -> &[String] {
self.inner.columns()
}
fn next(&mut self) -> bool {
self.inner.next()
}
fn scan(&self, dest: &mut [Value]) -> Result<()> {
self.inner.scan(dest)
}
fn row(&self) -> &Row {
self.inner.row()
}
fn take_row(&mut self) -> Row {
self.inner.take_row()
}
fn close(&mut self) -> Result<()> {
self.inner.close()
}
fn rows_affected(&self) -> i64 {
0
}
fn last_insert_id(&self) -> i64 {
0
}
fn with_aliases(self: Box<Self>, aliases: FxHashMap<String, String>) -> Box<dyn QueryResult> {
Box::new(AliasedResult::new(self, aliases))
}
}