use crate::spill::BufferedIndexedSpill;
use crate::{
PhysicalRow, RowSchema, RowSchemaExecution, ScalarEvalContext, ScalarOrder,
ScalarSubqueryRunner,
};
use std::cmp::Ordering;
use uqa_core::memory::ProductionControl;
use uqa_core::Value;
use uqa_sql::{SQLError, SQLParam, ScalarExpr};
const CACHED_ROWS: usize = 8;
struct CachedRow {
position: i64,
row: PhysicalRow,
order_key: Option<Vec<Value>>,
}
pub(super) struct PartitionRows<'a> {
rows: &'a mut BufferedIndexedSpill,
schema: RowSchema,
len: i64,
order_by: &'a [ScalarOrder],
params: &'a [SQLParam],
hook: &'a dyn uqa_sql::expr::EngineHook,
subqueries: &'a dyn ScalarSubqueryRunner,
cache: Vec<CachedRow>,
next_victim: usize,
}
impl<'a> PartitionRows<'a> {
pub(super) fn new(
rows: &'a mut BufferedIndexedSpill,
order_by: &'a [ScalarOrder],
params: &'a [SQLParam],
hook: &'a dyn uqa_sql::expr::EngineHook,
subqueries: &'a dyn ScalarSubqueryRunner,
) -> Result<Self, SQLError> {
let len = i64::try_from(rows.len())
.map_err(|_| SQLError::Internal("window partition row count exceeds i64".into()))?;
Ok(Self {
schema: rows.row_schema().clone(),
rows,
len,
order_by,
params,
hook,
subqueries,
cache: Vec::with_capacity(CACHED_ROWS),
next_victim: 0,
})
}
pub(super) const fn len(&self) -> i64 {
self.len
}
pub(super) const fn is_ordered(&self) -> bool {
!self.order_by.is_empty()
}
fn cached(&mut self, position: i64) -> Result<usize, SQLError> {
if let Some(index) = self
.cache
.iter()
.position(|entry| entry.position == position)
{
return Ok(index);
}
let index = u64::try_from(position).map_err(|_| out_of_partition(position))?;
let row = self.rows.get(index).map_err(super::exec_to_sql_error)?;
let entry = CachedRow {
position,
row,
order_key: None,
};
if self.cache.len() < CACHED_ROWS {
self.cache.push(entry);
Ok(self.cache.len() - 1)
} else {
let victim = self.next_victim;
self.next_victim = (victim + 1) % CACHED_ROWS;
self.cache[victim] = entry;
Ok(victim)
}
}
pub(super) fn row(&mut self, position: i64) -> Result<PhysicalRow, SQLError> {
let index = self.cached(position)?;
Ok(self.cache[index].row.clone())
}
pub(super) fn with_context<T>(
&mut self,
position: i64,
body: impl FnOnce(&ScalarEvalContext<'_>) -> Result<T, SQLError>,
) -> Result<T, SQLError> {
let row = self.row(position)?;
let view = self.schema.view(&row);
let context = ScalarEvalContext::from_row_lookup(&view, self.params)
.with_function_hook(self.hook)
.with_subquery_runner(self.subqueries)
.with_physical_outer_row(&self.schema, &row);
body(&context)
}
pub(super) fn enum_labels(&self) -> Option<&dyn uqa_sql::expr::enums::EnumLabelCatalog> {
self.hook.enum_labels()
}
pub(super) fn evaluate(
&mut self,
expression: &ScalarExpr,
position: i64,
) -> Result<Value, SQLError> {
let row = self.row(position)?;
super::evaluate_on_row(
expression,
&self.schema,
&row,
self.params,
self.hook,
self.subqueries,
)
}
fn order_key(&mut self, position: i64) -> Result<Vec<Value>, SQLError> {
let index = self.cached(position)?;
if let Some(key) = &self.cache[index].order_key {
return Ok(key.clone());
}
let row = self.cache[index].row.clone();
let key = self
.order_by
.iter()
.map(|order| {
super::evaluate_on_row(
&order.expr,
&self.schema,
&row,
self.params,
self.hook,
self.subqueries,
)
})
.collect::<Result<Vec<_>, _>>()?;
self.cache[index].order_key = Some(key.clone());
Ok(key)
}
pub(super) fn sort_value(&mut self, position: i64) -> Result<Value, SQLError> {
Ok(self
.order_key(position)?
.into_iter()
.next()
.unwrap_or(Value::Null))
}
pub(super) fn are_peers(&mut self, left: i64, right: i64) -> Result<bool, SQLError> {
if self.order_by.is_empty() || left == right {
return Ok(true);
}
let left = self.order_key(left)?;
let right = self.order_key(right)?;
for (left, right) in left.iter().zip(&right) {
let ordering = uqa_sql::expr::compare_typed_values_with_control(
left,
right,
&ProductionControl::uncontrolled(),
)?;
if ordering != Ordering::Equal {
return Ok(false);
}
}
Ok(true)
}
}
fn out_of_partition(position: i64) -> SQLError {
SQLError::Internal(format!("window row {position} lies outside its partition"))
}