use crate::{
db::{
data::{DecodedDataStoreKey, RawRow},
executor::{
ExecutionPreparation, PreparedGroupedRuntimeResidents,
aggregate::field::{
AggregateFieldValueError, FieldSlot,
extract_non_null_aggregate_field_value_with_slot_reader,
},
budget::{
charge_current_execution_budget, charge_decoded_row, charge_materialized_data_row,
},
pipeline::contracts::ResolvedExecutionKeyStream,
projection::{
eval_effective_runtime_filter_program_with_value_cow_reader,
resolve_slot_field_path, resolve_value_field_path,
},
terminal::{RetainedSlotLayout, RetainedSlotRow, RowDecoder, RowLayout},
},
predicate::MissingRowPolicy,
query::plan::{
EffectiveRuntimeFilterProgram, FieldSlot as PlannedFieldSlot, GroupFieldSet,
GroupedAggregateExecutionSpec, GroupedDistinctExecutionStrategy, ScalarGroupPath,
expr::{CompiledExprValueReader, ProjectionEvalError},
},
registry::StoreHandle,
},
error::InternalError,
value::Value,
};
use icydb_diagnostic_code::DiagnosticExecutionBudgetResource;
use std::{borrow::Cow, rc::Rc};
pub(in crate::db::executor) struct RowView {
storage: RowViewStorage,
}
pub(in crate::db::executor) fn compile_grouped_row_slot_layout_from_inputs(
row_layout: RowLayout,
group_fields: &GroupFieldSet,
grouped_aggregate_execution_specs: &[GroupedAggregateExecutionSpec],
grouped_distinct_execution_strategy: &GroupedDistinctExecutionStrategy,
effective_runtime_filter_program: Option<&EffectiveRuntimeFilterProgram>,
) -> RetainedSlotLayout {
let field_count = row_layout.field_count();
let mut required_slots = vec![false; field_count];
for field in group_fields.iter() {
if let Some(required_slot) = required_slots.get_mut(field.root_slot()) {
*required_slot = true;
}
}
if let Some(effective_runtime_filter_program) = effective_runtime_filter_program {
effective_runtime_filter_program.mark_referenced_slots(&mut required_slots);
}
for aggregate in grouped_aggregate_execution_specs {
if let Some(target_slot) = aggregate.target_slot()
&& let Some(required_slot) = required_slots.get_mut(target_slot.index())
{
*required_slot = true;
}
if let Some(compiled_input_expr) = aggregate.compiled_input_expr() {
compiled_input_expr.mark_referenced_slots(&mut required_slots);
}
if let Some(compiled_filter_expr) = aggregate.compiled_filter_expr() {
compiled_filter_expr.mark_referenced_slots(&mut required_slots);
}
}
if let Some(target_field) = grouped_distinct_execution_strategy.global_distinct_target_slot()
&& let Some(required_slot) = required_slots.get_mut(target_field.index())
{
*required_slot = true;
}
RetainedSlotLayout::compile(
field_count,
required_slots
.into_iter()
.enumerate()
.filter_map(|(slot, required)| required.then_some(slot))
.collect(),
)
}
enum RowViewStorage {
Single { slot: usize, value: Value },
SinglePath { value: Value },
Retained(RetainedSlotRow),
}
impl RowView {
#[must_use]
#[cfg(test)]
pub(in crate::db::executor) fn new(slots: Vec<Option<Value>>) -> Self {
let required_slots = slots
.iter()
.enumerate()
.filter_map(|(slot, value)| value.as_ref().map(|_| slot))
.collect();
let layout = RetainedSlotLayout::compile(slots.len(), required_slots);
let values = slots.into_iter().flatten().map(Some).collect();
Self::from_retained_slots(RetainedSlotRow::from_indexed_values(&layout, values))
}
#[must_use]
pub(in crate::db::executor) const fn from_retained_slots(row: RetainedSlotRow) -> Self {
Self {
storage: RowViewStorage::Retained(row),
}
}
#[must_use]
pub(in crate::db::executor) const fn from_single_value(slot: usize, value: Value) -> Self {
Self {
storage: RowViewStorage::Single { slot, value },
}
}
pub(in crate::db::executor) fn slot_value_ref(&self, index: usize) -> Option<&Value> {
match &self.storage {
RowViewStorage::Single { slot, value } => (*slot == index).then_some(value),
RowViewStorage::SinglePath { .. } => None,
RowViewStorage::Retained(row) => row.slot_ref(index),
}
}
pub(in crate::db::executor) fn require_slot_value(
&self,
index: usize,
) -> Result<&Value, InternalError> {
self.slot_value_ref(index)
.ok_or_else(InternalError::query_executor_invariant)
}
pub(in crate::db::executor) fn into_required_slot_value(
self,
index: usize,
) -> Result<Value, InternalError> {
match self.storage {
RowViewStorage::Single { slot, value } => {
if slot == index {
return Ok(value);
}
Err(InternalError::query_executor_invariant())
}
RowViewStorage::SinglePath { .. } => Err(InternalError::query_executor_invariant()),
RowViewStorage::Retained(mut row) => row
.take_slot(index)
.ok_or_else(InternalError::query_executor_invariant),
}
}
pub(in crate::db::executor) fn require_slot_owned(
&self,
index: usize,
) -> Result<Value, InternalError> {
self.require_slot_value(index).cloned()
}
pub(in crate::db::executor) fn with_required_slot<R>(
&self,
index: usize,
f: impl FnOnce(&Value) -> Result<R, InternalError>,
) -> Result<R, InternalError> {
f(self.require_slot_value(index)?)
}
#[must_use]
pub(in crate::db::executor) const fn predecoded_single_group_path_value(
&self,
) -> Option<&Value> {
match &self.storage {
RowViewStorage::SinglePath { value } => Some(value),
RowViewStorage::Single { .. } | RowViewStorage::Retained(_) => None,
}
}
pub(in crate::db::executor) fn eval_filter_program(
&self,
effective_runtime_filter_program: &EffectiveRuntimeFilterProgram,
) -> Result<bool, InternalError> {
eval_effective_runtime_filter_program_with_value_cow_reader(
effective_runtime_filter_program,
&mut |slot| self.slot_value_ref(slot).map(Cow::Borrowed),
)
}
pub(in crate::db::executor) fn extract_non_null_aggregate_field_value(
&self,
field_slot: FieldSlot,
) -> Result<Option<Value>, InternalError> {
let mut value = Some(self.require_slot_owned(field_slot.index)?);
extract_non_null_aggregate_field_value_with_slot_reader(field_slot, &mut |_| value.take())
.map_err(AggregateFieldValueError::into_internal_error)
}
pub(in crate::db::executor) fn group_values(
&self,
group_fields: &[PlannedFieldSlot],
) -> Result<Vec<Value>, InternalError> {
let mut values = Vec::with_capacity(group_fields.len());
for field in group_fields {
let value = self.require_slot_owned(field.index())?;
values.push(value);
}
Ok(values)
}
}
impl CompiledExprValueReader for RowView {
fn read_slot(&self, slot: usize) -> Option<Cow<'_, Value>> {
self.slot_value_ref(slot).map(Cow::Borrowed)
}
fn read_group_key(&self, _offset: usize) -> Option<Cow<'_, Value>> {
None
}
fn read_aggregate(&self, _index: usize) -> Option<Cow<'_, Value>> {
None
}
fn read_field_path(
&self,
root_slot: usize,
segments: &[String],
_segment_bytes: &[Box<[u8]>],
) -> Result<Option<Cow<'_, Value>>, ProjectionEvalError> {
let Some(root) = self.slot_value_ref(root_slot) else {
return Ok(None);
};
let value = resolve_value_field_path(root, segments)?;
Ok(Some(value.map_or(Cow::Owned(Value::Null), Cow::Borrowed)))
}
}
struct SingleGroupedSlotDecode {
slot: usize,
}
struct SingleGroupedPathDecode {
root_slot: usize,
segment_bytes: Box<[Box<[u8]>]>,
}
impl SingleGroupedPathDecode {
fn new(path: &ScalarGroupPath) -> Self {
Self {
root_slot: path.root_slot(),
segment_bytes: path
.path()
.segments()
.iter()
.map(|segment| segment.as_bytes().to_vec().into_boxed_slice())
.collect::<Vec<_>>()
.into_boxed_slice(),
}
}
}
enum SingleGroupedDecode {
Path(SingleGroupedPathDecode),
Slot(SingleGroupedSlotDecode),
}
pub(in crate::db::executor) struct StructuralGroupedRowRuntime {
store: StoreHandle,
row_layout: RowLayout,
grouped_slot_layout: RetainedSlotLayout,
single_grouped_decode: Option<SingleGroupedDecode>,
}
impl StructuralGroupedRowRuntime {
#[must_use]
pub(in crate::db::executor) fn new(
store: StoreHandle,
row_layout: RowLayout,
grouped_slot_layout: RetainedSlotLayout,
single_grouped_path: Option<&ScalarGroupPath>,
) -> Self {
let single_grouped_decode = single_grouped_path
.map(SingleGroupedPathDecode::new)
.map(SingleGroupedDecode::Path)
.or_else(|| match grouped_slot_layout.required_slots() {
[required_slot] => Some(SingleGroupedDecode::Slot(SingleGroupedSlotDecode {
slot: *required_slot,
})),
_ => None,
});
Self {
store,
row_layout,
grouped_slot_layout,
single_grouped_decode,
}
}
fn row_view_from_data_row(
&self,
key: &DecodedDataStoreKey,
row: RawRow,
) -> Result<RowView, InternalError> {
match self.single_grouped_decode.as_ref() {
Some(SingleGroupedDecode::Path(single_grouped_path_decode)) => {
self.single_path_row_view_from_data_row(key, row, single_grouped_path_decode)
}
Some(SingleGroupedDecode::Slot(single_grouped_slot_decode)) => {
self.single_slot_row_view_from_data_row(key, row, single_grouped_slot_decode)
}
None => {
charge_decoded_row(row.len(), self.grouped_slot_layout.required_slots().len())?;
let retained_slots = RowDecoder::decode_retained_slots_from_data_key(
&self.row_layout,
key,
&row,
&self.grouped_slot_layout,
)?;
Ok(RowView::from_retained_slots(retained_slots))
}
}
}
fn single_path_row_view_from_data_row(
&self,
key: &DecodedDataStoreKey,
row: RawRow,
path: &SingleGroupedPathDecode,
) -> Result<RowView, InternalError> {
charge_decoded_row(row.len(), path.segment_bytes.len())?;
let row_fields = self.row_layout.open_raw_row_with_contract(&row)?;
row_fields.validate_primary_key(key)?;
let value = resolve_slot_field_path(&row_fields, path.root_slot, &path.segment_bytes)?
.unwrap_or(Value::Null);
Ok(RowView {
storage: RowViewStorage::SinglePath { value },
})
}
fn single_slot_row_view_from_data_row(
&self,
key: &DecodedDataStoreKey,
row: RawRow,
single_grouped_slot_decode: &SingleGroupedSlotDecode,
) -> Result<RowView, InternalError> {
let value = self.decode_single_grouped_slot_value_from_raw_row(
key,
&row,
single_grouped_slot_decode,
)?;
Ok(RowView::from_single_value(
single_grouped_slot_decode.slot,
value,
))
}
fn decode_single_grouped_slot_value_from_raw_row(
&self,
key: &DecodedDataStoreKey,
row: &RawRow,
single_grouped_slot_decode: &SingleGroupedSlotDecode,
) -> Result<Value, InternalError> {
charge_decoded_row(row.len(), 1)?;
RowLayout::decode_required_value_from_data_key(
&self.row_layout,
row,
key,
single_grouped_slot_decode.slot,
)
}
const fn matching_single_grouped_slot_decode(
&self,
required_slot: usize,
) -> Option<&SingleGroupedSlotDecode> {
match self.single_grouped_decode.as_ref() {
Some(SingleGroupedDecode::Slot(single_grouped_slot_decode))
if single_grouped_slot_decode.slot == required_slot =>
{
Some(single_grouped_slot_decode)
}
Some(SingleGroupedDecode::Path(_) | SingleGroupedDecode::Slot(_)) | None => None,
}
}
fn read_data_row(
&self,
consistency: MissingRowPolicy,
key: &DecodedDataStoreKey,
) -> Result<Option<RawRow>, InternalError> {
let raw_key = key.to_raw()?;
let row = self.store.with_data(|store| store.get(&raw_key));
charge_current_execution_budget(DiagnosticExecutionBudgetResource::RowsVisited, 1)?;
if let Some(row) = row.as_ref() {
charge_materialized_data_row!(row)?;
}
match (consistency, row) {
(MissingRowPolicy::Ignore, None) => Ok(None),
(MissingRowPolicy::Ignore | MissingRowPolicy::Error, Some(row)) => Ok(Some(row)),
(MissingRowPolicy::Error, None) => Err(InternalError::store_corruption()),
}
}
pub(in crate::db::executor) fn read_single_group_value(
&self,
consistency: MissingRowPolicy,
key: &DecodedDataStoreKey,
required_slot: usize,
) -> Result<Option<Value>, InternalError> {
let Some(row) = self.read_data_row(consistency, key)? else {
return Ok(None);
};
if let Some(single_grouped_slot_decode) =
self.matching_single_grouped_slot_decode(required_slot)
{
return self
.decode_single_grouped_slot_value_from_raw_row(
key,
&row,
single_grouped_slot_decode,
)
.map(Some);
}
let row_view = self.row_view_from_data_row(key, row)?;
row_view.into_required_slot_value(required_slot).map(Some)
}
pub(in crate::db::executor) fn read_row_view(
&self,
consistency: MissingRowPolicy,
key: &DecodedDataStoreKey,
) -> Result<Option<RowView>, InternalError> {
self.read_data_row(consistency, key)?
.map(|row| self.row_view_from_data_row(key, row))
.transpose()
}
}
pub(in crate::db::executor) struct GroupedStreamStage {
row_runtime: StructuralGroupedRowRuntime,
prepared_residents: Rc<PreparedGroupedRuntimeResidents>,
resolved: ResolvedExecutionKeyStream,
}
impl GroupedStreamStage {
pub(in crate::db::executor) const fn new(
row_runtime: StructuralGroupedRowRuntime,
prepared_residents: Rc<PreparedGroupedRuntimeResidents>,
resolved: ResolvedExecutionKeyStream,
) -> Self {
Self {
row_runtime,
prepared_residents,
resolved,
}
}
pub(in crate::db::executor) fn fold_inputs_mut(
&mut self,
) -> (
&StructuralGroupedRowRuntime,
&ExecutionPreparation,
&mut ResolvedExecutionKeyStream,
) {
(
&self.row_runtime,
self.prepared_residents.execution_preparation(),
&mut self.resolved,
)
}
}
#[cfg(test)]
mod tests {
use crate::{
db::executor::{
pipeline::runtime::RowView,
terminal::{RetainedSlotLayout, RetainedSlotRow},
},
value::Value,
};
#[test]
fn retained_row_view_resolves_noncontiguous_slots() {
let row_view = RowView::new(vec![
None,
Some(Value::Nat64(7)),
None,
None,
Some(Value::Text("group".to_string())),
None,
]);
assert_eq!(row_view.slot_value_ref(1), Some(&Value::Nat64(7)));
assert_eq!(
row_view.slot_value_ref(4),
Some(&Value::Text("group".to_string()))
);
assert_eq!(row_view.slot_value_ref(0), None);
assert_eq!(row_view.slot_value_ref(5), None);
assert_eq!(row_view.slot_value_ref(6), None);
assert_eq!(
row_view.into_required_slot_value(4).unwrap(),
Value::Text("group".to_string())
);
}
#[test]
fn single_slot_row_view_resolves_only_its_declared_slot() {
let row_view = RowView::from_single_value(4, Value::Text("group".to_string()));
assert_eq!(
row_view.slot_value_ref(4),
Some(&Value::Text("group".to_string()))
);
assert_eq!(row_view.slot_value_ref(1), None);
}
#[test]
fn retained_row_view_slot_reads_are_repeatable_borrows() {
let layout = RetainedSlotLayout::compile(5, vec![1, 4]);
let retained = RetainedSlotRow::from_indexed_values(
&layout,
vec![
Some(Value::Nat64(7)),
Some(Value::Text("group".to_string())),
],
);
let row_view = RowView::from_retained_slots(retained);
assert_eq!(
row_view.slot_value_ref(1),
Some(&Value::Nat64(7)),
"first retained slot read should borrow the decoded value",
);
assert_eq!(
row_view.slot_value_ref(1),
Some(&Value::Nat64(7)),
"second retained slot read must see the same decoded value",
);
assert_eq!(
row_view.slot_value_ref(4),
Some(&Value::Text("group".to_string())),
"reading another retained slot must not invalidate earlier slots",
);
}
}