use std::fmt;
use radixdb_core::value::NULL_VALUE;
use radixdb_core::{CompactArc, CompactVec};
use radixdb_core::{Error, Result, Row, Value};
use radixdb_storage::{DeferredColumnSource, DeferredRow};
#[derive(Debug, Clone)]
pub struct ColumnInfo {
pub name: String,
pub table_alias: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub enum OrderingProperty {
#[default]
Unknown,
AscendingNullsLast(Vec<usize>),
}
impl OrderingProperty {
pub fn ascending_nulls_last(key_indices: Vec<usize>) -> Self {
if key_indices.is_empty() {
Self::Unknown
} else {
Self::AscendingNullsLast(key_indices)
}
}
pub fn proves_ascending_nulls_last(&self, required_keys: &[usize]) -> bool {
if required_keys.is_empty() {
return false;
}
match self {
Self::AscendingNullsLast(keys) => keys.starts_with(required_keys),
Self::Unknown => false,
}
}
pub fn remap_outer_projection(&self, columns: &[ColumnSource]) -> Self {
let Self::AscendingNullsLast(keys) = self else {
return Self::Unknown;
};
let mut remapped = Vec::with_capacity(keys.len());
for key in keys {
let Some(output) = columns
.iter()
.position(|source| matches!(source, ColumnSource::Outer(index) if index == key))
else {
return Self::Unknown;
};
remapped.push(output);
}
Self::ascending_nulls_last(remapped)
}
}
impl ColumnInfo {
pub fn new(name: impl Into<String>) -> Self {
Self {
name: name.into(),
table_alias: None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ColumnSource {
Outer(usize),
Inner(usize),
}
#[derive(Debug, Clone)]
pub struct JoinProjection {
pub columns: Vec<ColumnSource>,
}
impl JoinProjection {
pub fn validate(
&self,
left_columns: usize,
right_columns: usize,
output_columns: usize,
) -> Result<()> {
if self.columns.len() != output_columns {
return Err(Error::invalid_argument(format!(
"join projection has {} values but output schema has {} columns",
self.columns.len(),
output_columns
)));
}
for source in &self.columns {
let (side, index, width) = match source {
ColumnSource::Outer(index) => ("left", *index, left_columns),
ColumnSource::Inner(index) => ("right", *index, right_columns),
};
if index >= width {
return Err(Error::invalid_argument(format!(
"join projection {side} column index {index} is outside width {width}"
)));
}
}
Ok(())
}
}
pub trait Operator: Send {
fn open(&mut self) -> Result<()>;
fn next(&mut self) -> Result<Option<RowRef>>;
fn close(&mut self) -> Result<()>;
fn schema(&self) -> &[ColumnInfo];
fn estimated_rows(&self) -> Option<usize> {
None
}
fn ordering(&self) -> OrderingProperty {
OrderingProperty::Unknown
}
fn name(&self) -> &str;
}
#[derive(Debug, Clone)]
pub enum RowRef {
Owned(Row),
Composite(CompositeRow),
DirectBuildComposite(DirectBuildCompositeRow),
Shared(SharedRow),
Projected(ProjectedRow),
Deferred(DeferredRow),
}
impl RowRef {
#[inline]
pub fn owned(row: Row) -> Self {
RowRef::Owned(row)
}
#[inline]
pub fn composite(left: Row, right: Row) -> Self {
RowRef::Composite(CompositeRow::new(left, right))
}
#[inline]
pub fn len(&self) -> usize {
match self {
RowRef::Owned(row) => row.len(),
RowRef::Composite(comp) => comp.len(),
RowRef::DirectBuildComposite(direct) => direct.len(),
RowRef::Shared(shared) => shared.len(),
RowRef::Projected(projected) => projected.len(),
RowRef::Deferred(deferred) => deferred.len(),
}
}
#[inline]
pub fn is_empty(&self) -> bool {
self.len() == 0
}
#[inline]
pub fn get(&self, idx: usize) -> Option<&Value> {
match self {
RowRef::Owned(row) => row.get(idx),
RowRef::Composite(comp) => comp.get(idx),
RowRef::DirectBuildComposite(direct) => direct.get(idx),
RowRef::Shared(shared) => shared.get(idx),
RowRef::Projected(projected) => projected.get(idx),
RowRef::Deferred(deferred) => deferred.get(idx),
}
}
#[inline]
pub fn into_owned(self) -> Row {
match self {
RowRef::Owned(row) => row,
RowRef::Composite(comp) => comp.materialize_owned(),
RowRef::DirectBuildComposite(direct) => direct.materialize_owned(),
RowRef::Shared(shared) => shared.materialize_owned(),
RowRef::Projected(projected) => projected.materialize_owned(),
RowRef::Deferred(deferred) => deferred.into_owned(),
}
}
pub fn to_owned(&self) -> Row {
match self {
RowRef::Owned(row) => row.clone(),
RowRef::Composite(comp) => comp.materialize(),
RowRef::DirectBuildComposite(direct) => direct.materialize(),
RowRef::Shared(shared) => shared.materialize(),
RowRef::Projected(projected) => projected.materialize(),
RowRef::Deferred(deferred) => deferred.to_owned(),
}
}
#[inline]
pub fn as_row(&self) -> Option<&Row> {
match self {
RowRef::Owned(row) => Some(row),
RowRef::Shared(shared) => Some(shared.row()),
RowRef::Composite(_)
| RowRef::DirectBuildComposite(_)
| RowRef::Projected(_)
| RowRef::Deferred(_) => None,
}
}
#[inline]
pub fn direct_build_composite(
probe: Row,
build_rows: CompactArc<Vec<Row>>,
build_idx: usize,
probe_is_left: bool,
) -> Self {
RowRef::DirectBuildComposite(DirectBuildCompositeRow::new(
probe,
build_rows,
build_idx,
probe_is_left,
))
}
#[inline]
pub fn shared(rows: CompactArc<Vec<Row>>, row_idx: usize) -> Self {
RowRef::Shared(SharedRow::new(rows, row_idx))
}
#[inline]
pub fn projected(left: RowRef, right: RowRef, columns: CompactArc<[ColumnSource]>) -> Self {
RowRef::Projected(ProjectedRow::new(left, right, columns))
}
#[inline]
pub fn deferred(row: DeferredRow) -> Self {
RowRef::Deferred(row)
}
pub fn into_deferred(self) -> DeferredRow {
match self {
RowRef::Owned(row) => DeferredRow::owned(row),
RowRef::Shared(SharedRow { rows, row_idx }) => DeferredRow::shared(rows, row_idx),
RowRef::Projected(ProjectedRow {
left,
right,
columns,
}) => {
let columns = columns
.iter()
.map(|source| match source {
ColumnSource::Outer(index) => DeferredColumnSource::Left(*index),
ColumnSource::Inner(index) => DeferredColumnSource::Right(*index),
})
.collect::<Vec<_>>();
DeferredRow::projected(
left.into_deferred(),
right.into_deferred(),
CompactArc::from(columns),
)
}
RowRef::Composite(composite) => DeferredRow::owned(composite.materialize_owned()),
RowRef::DirectBuildComposite(composite) => {
DeferredRow::owned(composite.materialize_owned())
}
RowRef::Deferred(row) => row,
}
}
#[inline]
pub fn is_deferred(&self) -> bool {
match self {
RowRef::Owned(_) => false,
RowRef::Deferred(row) => row.is_deferred(),
_ => true,
}
}
pub fn estimated_retained_bytes(&self) -> usize {
fn row_bytes(row: &Row) -> usize {
row.iter().fold(std::mem::size_of::<Row>(), |total, value| {
let payload = match value {
Value::Text(text) => text.len(),
Value::Extension(bytes) => bytes.len(),
_ => 0,
};
total
.saturating_add(std::mem::size_of::<Value>())
.saturating_add(payload)
})
}
match self {
Self::Owned(row) => row_bytes(row),
Self::Shared(_) => std::mem::size_of::<Self>(),
Self::Composite(row) => std::mem::size_of::<Self>()
.saturating_add(row_bytes(&row.left))
.saturating_add(row_bytes(&row.right)),
Self::DirectBuildComposite(row) => {
std::mem::size_of::<Self>().saturating_add(row_bytes(&row.probe))
}
Self::Projected(row) => std::mem::size_of::<Self>()
.saturating_add(row.left.estimated_retained_bytes())
.saturating_add(row.right.estimated_retained_bytes())
.saturating_add(
row.columns
.len()
.saturating_mul(std::mem::size_of::<ColumnSource>()),
),
Self::Deferred(row) => row.estimated_retained_bytes(),
}
}
}
#[derive(Debug, Clone)]
pub struct SharedRow {
rows: CompactArc<Vec<Row>>,
row_idx: usize,
}
impl SharedRow {
#[inline]
pub fn new(rows: CompactArc<Vec<Row>>, row_idx: usize) -> Self {
assert!(
row_idx < rows.len(),
"shared row index outside immutable batch"
);
Self { rows, row_idx }
}
#[inline]
fn row(&self) -> &Row {
&self.rows[self.row_idx]
}
#[inline]
pub fn len(&self) -> usize {
self.row().len()
}
#[inline]
pub fn is_empty(&self) -> bool {
self.row().is_empty()
}
#[inline]
pub fn get(&self, idx: usize) -> Option<&Value> {
self.row().get(idx)
}
#[inline]
pub fn materialize(&self) -> Row {
self.row().clone()
}
#[inline]
pub fn materialize_owned(self) -> Row {
self.row().clone()
}
}
#[derive(Debug, Clone)]
pub struct ProjectedRow {
left: Box<RowRef>,
right: Box<RowRef>,
columns: CompactArc<[ColumnSource]>,
}
impl ProjectedRow {
#[inline]
pub fn new(left: RowRef, right: RowRef, columns: CompactArc<[ColumnSource]>) -> Self {
Self {
left: Box::new(left),
right: Box::new(right),
columns,
}
}
#[inline]
pub fn len(&self) -> usize {
self.columns.len()
}
#[inline]
pub fn is_empty(&self) -> bool {
self.columns.is_empty()
}
#[inline]
pub fn get(&self, idx: usize) -> Option<&Value> {
match self.columns.get(idx)? {
ColumnSource::Outer(index) => self.left.get(*index),
ColumnSource::Inner(index) => self.right.get(*index),
}
}
pub fn materialize(&self) -> Row {
let mut values = CompactVec::with_capacity(self.columns.len());
for index in 0..self.columns.len() {
values.push(self.get(index).cloned().unwrap_or(NULL_VALUE));
}
Row::from_compact_vec(values)
}
#[inline]
pub fn materialize_owned(self) -> Row {
self.materialize()
}
}
impl fmt::Display for ProjectedRow {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "(")?;
for index in 0..self.len() {
if index > 0 {
write!(f, ", ")?;
}
match self.get(index) {
Some(value) => write!(f, "{value}")?,
None => write!(f, "NULL")?,
}
}
write!(f, ")")
}
}
#[derive(Debug, Clone)]
pub struct CompositeRow {
left: Row,
right: Row,
left_cols: usize,
}
impl CompositeRow {
#[inline]
pub fn new(left: Row, right: Row) -> Self {
let left_cols = left.len();
Self {
left,
right,
left_cols,
}
}
#[inline]
pub fn len(&self) -> usize {
self.left_cols + self.right.len()
}
#[inline]
pub fn is_empty(&self) -> bool {
self.left.is_empty() && self.right.is_empty()
}
#[inline]
pub fn get(&self, idx: usize) -> Option<&Value> {
if idx < self.left_cols {
self.left.get(idx)
} else {
self.right.get(idx - self.left_cols)
}
}
#[inline]
pub fn left(&self) -> &Row {
&self.left
}
#[inline]
pub fn right(&self) -> &Row {
&self.right
}
pub fn materialize(&self) -> Row {
let total = self.len();
let mut values: CompactVec<Value> = CompactVec::with_capacity(total);
values.extend_clone(self.left.as_slice());
values.extend_clone(self.right.as_slice());
Row::from_compact_vec(values)
}
#[inline]
pub fn materialize_owned(self) -> Row {
Row::from_combined_owned(self.left, self.right)
}
#[inline]
pub fn into_parts(self) -> (Row, Row) {
(self.left, self.right)
}
}
impl fmt::Display for CompositeRow {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "(")?;
for i in 0..self.len() {
if i > 0 {
write!(f, ", ")?;
}
if let Some(v) = self.get(i) {
write!(f, "{}", v)?;
} else {
write!(f, "NULL")?;
}
}
write!(f, ")")
}
}
#[derive(Debug)]
pub struct DirectBuildCompositeRow {
probe: Row,
build_rows: CompactArc<Vec<Row>>,
build_idx: usize,
probe_cols: usize,
probe_is_left: bool,
}
impl DirectBuildCompositeRow {
#[inline]
pub fn new(
probe: Row,
build_rows: CompactArc<Vec<Row>>,
build_idx: usize,
probe_is_left: bool,
) -> Self {
debug_assert!(
build_idx < build_rows.len(),
"build_idx {} out of bounds (len={})",
build_idx,
build_rows.len()
);
let probe_cols = probe.len();
Self {
probe,
build_rows,
build_idx,
probe_cols,
probe_is_left,
}
}
#[inline]
pub fn len(&self) -> usize {
self.probe_cols + self.build_rows[self.build_idx].len()
}
#[inline]
pub fn is_empty(&self) -> bool {
self.probe.is_empty() && self.build_rows[self.build_idx].is_empty()
}
#[inline]
pub fn get(&self, idx: usize) -> Option<&Value> {
let build_row = &self.build_rows[self.build_idx];
if self.probe_is_left {
if idx < self.probe_cols {
self.probe.get(idx)
} else {
build_row.get(idx - self.probe_cols)
}
} else {
let build_cols = build_row.len();
if idx < build_cols {
build_row.get(idx)
} else {
self.probe.get(idx - build_cols)
}
}
}
pub fn materialize(&self) -> Row {
let build_row = &self.build_rows[self.build_idx];
let total = self.probe_cols + build_row.len();
let mut values: CompactVec<Value> = CompactVec::with_capacity(total);
if self.probe_is_left {
values.extend_clone(self.probe.as_slice());
values.extend_clone(build_row.as_slice());
} else {
values.extend_clone(build_row.as_slice());
values.extend_clone(self.probe.as_slice());
}
Row::from_compact_vec(values)
}
#[inline]
pub fn materialize_owned(self) -> Row {
let build_row = &self.build_rows[self.build_idx];
let total = self.probe_cols + build_row.len();
let mut values: CompactVec<Value> = CompactVec::with_capacity(total);
if self.probe_is_left {
self.probe.extend_into_compact_vec(&mut values);
values.extend_clone(build_row.as_slice());
} else {
values.extend_clone(build_row.as_slice());
self.probe.extend_into_compact_vec(&mut values);
}
Row::from_compact_vec(values)
}
}
impl Clone for DirectBuildCompositeRow {
fn clone(&self) -> Self {
Self {
probe: self.probe.clone(),
build_rows: CompactArc::clone(&self.build_rows),
build_idx: self.build_idx,
probe_cols: self.probe_cols,
probe_is_left: self.probe_is_left,
}
}
}
impl fmt::Display for DirectBuildCompositeRow {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "(")?;
for i in 0..self.len() {
if i > 0 {
write!(f, ", ")?;
}
if let Some(v) = self.get(i) {
write!(f, "{}", v)?;
} else {
write!(f, "NULL")?;
}
}
write!(f, ")")
}
}
pub struct EmptyOperator {
schema: Vec<ColumnInfo>,
opened: bool,
}
impl EmptyOperator {
pub fn new() -> Self {
Self {
schema: Vec::new(),
opened: false,
}
}
}
impl Default for EmptyOperator {
fn default() -> Self {
Self::new()
}
}
impl Operator for EmptyOperator {
fn open(&mut self) -> Result<()> {
self.opened = true;
Ok(())
}
fn next(&mut self) -> Result<Option<RowRef>> {
Ok(None)
}
fn close(&mut self) -> Result<()> {
Ok(())
}
fn schema(&self) -> &[ColumnInfo] {
&self.schema
}
fn name(&self) -> &str {
"Empty"
}
}
pub struct MaterializedOperator {
rows: Vec<Row>,
schema: Vec<ColumnInfo>,
ordering: OrderingProperty,
current_idx: usize,
opened: bool,
}
impl MaterializedOperator {
pub fn new(rows: Vec<Row>, schema: Vec<ColumnInfo>) -> Self {
Self {
rows,
schema,
ordering: OrderingProperty::Unknown,
current_idx: 0,
opened: false,
}
}
pub fn with_ordering(mut self, ordering: OrderingProperty) -> Self {
self.ordering = ordering;
self
}
pub fn from_arc(arc_rows: CompactArc<Vec<Row>>, schema: Vec<ColumnInfo>) -> Self {
let rows = CompactArc::try_unwrap(arc_rows).unwrap_or_else(|arc| (*arc).clone());
Self::new(rows, schema)
}
}
impl Operator for MaterializedOperator {
fn open(&mut self) -> Result<()> {
self.current_idx = 0;
self.opened = true;
Ok(())
}
fn next(&mut self) -> Result<Option<RowRef>> {
if self.current_idx >= self.rows.len() {
return Ok(None);
}
let row = std::mem::take(&mut self.rows[self.current_idx]);
self.current_idx += 1;
Ok(Some(RowRef::Owned(row)))
}
fn close(&mut self) -> Result<()> {
Ok(())
}
fn schema(&self) -> &[ColumnInfo] {
&self.schema
}
fn estimated_rows(&self) -> Option<usize> {
Some(self.rows.len())
}
fn ordering(&self) -> OrderingProperty {
self.ordering.clone()
}
fn name(&self) -> &str {
"Materialized"
}
}
use radixdb_storage::QueryResult as StorageQueryResult;
pub struct QueryResultOperator {
result: Box<dyn StorageQueryResult>,
schema: Vec<ColumnInfo>,
ordering: OrderingProperty,
opened: bool,
}
impl QueryResultOperator {
pub fn new(result: Box<dyn StorageQueryResult>, columns: Vec<String>) -> Self {
let ordering = result.ascending_nulls_last_ordering().map_or(
OrderingProperty::Unknown,
OrderingProperty::ascending_nulls_last,
);
let schema = columns.into_iter().map(ColumnInfo::new).collect();
Self {
result,
schema,
ordering,
opened: false,
}
}
}
impl Operator for QueryResultOperator {
fn open(&mut self) -> Result<()> {
radixdb_storage::instrumentation::record_join_source_open(self.opened);
self.opened = true;
Ok(())
}
fn next(&mut self) -> Result<Option<RowRef>> {
if !self.opened {
return Ok(None);
}
if self.result.next() {
Ok(Some(RowRef::deferred(self.result.take_deferred_row())))
} else if let Some(error) = self.result.last_error() {
Err(error)
} else {
Ok(None)
}
}
fn close(&mut self) -> Result<()> {
self.result.close()
}
fn schema(&self) -> &[ColumnInfo] {
&self.schema
}
fn estimated_rows(&self) -> Option<usize> {
None
}
fn ordering(&self) -> OrderingProperty {
self.ordering.clone()
}
fn name(&self) -> &str {
"QueryResultScan"
}
}
#[cfg(test)]
#[allow(clippy::approx_constant)]
mod tests {
use super::*;
#[test]
fn join_projection_rejects_shape_and_side_indices_before_execution() {
let wrong_width = JoinProjection {
columns: vec![ColumnSource::Outer(0)],
};
assert!(wrong_width.validate(1, 1, 2).is_err());
let bad_left = JoinProjection {
columns: vec![ColumnSource::Outer(1)],
};
assert!(bad_left.validate(1, 1, 1).is_err());
let bad_right = JoinProjection {
columns: vec![ColumnSource::Inner(1)],
};
assert!(bad_right.validate(1, 1, 1).is_err());
let valid = JoinProjection {
columns: vec![ColumnSource::Inner(0), ColumnSource::Outer(0)],
};
valid.validate(1, 1, 2).unwrap();
}
#[test]
fn test_composite_row_basic() {
let left = Row::from_values(vec![Value::integer(1), Value::text("hello")]);
let right = Row::from_values(vec![Value::float(3.14), Value::boolean(true)]);
let comp = CompositeRow::new(left, right);
assert_eq!(comp.len(), 4);
assert_eq!(comp.get(0), Some(&Value::integer(1)));
assert_eq!(comp.get(1), Some(&Value::text("hello")));
assert_eq!(comp.get(2), Some(&Value::float(3.14)));
assert_eq!(comp.get(3), Some(&Value::boolean(true)));
assert_eq!(comp.get(4), None);
}
#[test]
fn test_composite_row_materialize() {
let left = Row::from_values(vec![Value::integer(1)]);
let right = Row::from_values(vec![Value::integer(2)]);
let comp = CompositeRow::new(left, right);
let materialized = comp.materialize();
assert_eq!(materialized.len(), 2);
assert_eq!(materialized.get(0), Some(&Value::integer(1)));
assert_eq!(materialized.get(1), Some(&Value::integer(2)));
}
#[test]
fn test_row_ref_owned() {
let row = Row::from_values(vec![Value::integer(42)]);
let row_ref = RowRef::owned(row);
assert_eq!(row_ref.len(), 1);
assert_eq!(row_ref.get(0), Some(&Value::integer(42)));
let owned = row_ref.into_owned();
assert_eq!(owned.get(0), Some(&Value::integer(42)));
}
#[test]
fn test_row_ref_composite() {
let left = Row::from_values(vec![Value::integer(1)]);
let right = Row::from_values(vec![Value::integer(2)]);
let row_ref = RowRef::composite(left, right);
assert_eq!(row_ref.len(), 2);
assert_eq!(row_ref.get(0), Some(&Value::integer(1)));
assert_eq!(row_ref.get(1), Some(&Value::integer(2)));
}
#[test]
fn projected_row_ref_keeps_transitive_join_slots_deferred() {
let first = RowRef::projected(
RowRef::owned(Row::from_values(vec![
Value::integer(1),
Value::text("payload"),
])),
RowRef::owned(Row::from_values(vec![
Value::integer(10),
Value::text("dictionary"),
])),
CompactArc::from(vec![
ColumnSource::Outer(0),
ColumnSource::Outer(1),
ColumnSource::Inner(1),
]),
);
let second = RowRef::projected(
first,
RowRef::owned(Row::from_values(vec![
Value::integer(20),
Value::text("leaf"),
])),
CompactArc::from(vec![
ColumnSource::Outer(1),
ColumnSource::Inner(1),
ColumnSource::Outer(2),
]),
);
assert!(second.is_deferred());
assert_eq!(second.len(), 3);
assert_eq!(second.get(0), Some(&Value::text("payload")));
assert_eq!(second.get(1), Some(&Value::text("leaf")));
assert_eq!(second.get(2), Some(&Value::text("dictionary")));
let materialized = second.into_owned();
assert_eq!(
materialized,
Row::from_values(vec![
Value::text("payload"),
Value::text("leaf"),
Value::text("dictionary"),
])
);
}
#[test]
fn test_empty_operator() {
let mut op = EmptyOperator::new();
op.open().unwrap();
assert!(op.next().unwrap().is_none());
assert!(op.next().unwrap().is_none());
op.close().unwrap();
}
#[test]
fn test_materialized_operator() {
let rows = vec![
Row::from_values(vec![Value::integer(1)]),
Row::from_values(vec![Value::integer(2)]),
Row::from_values(vec![Value::integer(3)]),
];
let schema = vec![ColumnInfo::new("id")];
let mut op = MaterializedOperator::new(rows, schema);
op.open().unwrap();
let row1 = op.next().unwrap().unwrap();
assert_eq!(row1.get(0), Some(&Value::integer(1)));
let row2 = op.next().unwrap().unwrap();
assert_eq!(row2.get(0), Some(&Value::integer(2)));
let row3 = op.next().unwrap().unwrap();
assert_eq!(row3.get(0), Some(&Value::integer(3)));
assert!(op.next().unwrap().is_none());
op.close().unwrap();
}
}