use crate::{
db::{
data::{
CanonicalSlotReader, DecodedDataStoreKey, RawRow, decode_structural_value_storage_bytes,
},
executor::{
ExecutionPreparation, PreparedGroupedRuntimeResidents,
aggregate::field::{
AggregateFieldValueError, FieldSlot, extract_orderable_field_value_with_slot_reader,
},
budget::{charge_current_execution_budget, charge_materialized_data_row},
pipeline::contracts::ResolvedExecutionKeyStream,
projection::{
eval_effective_runtime_filter_program_with_value_cow_reader, resolve_path_segments,
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 {
#[cfg(test)]
Dense(Vec<Option<Value>>),
Single {
slot: usize,
value: Value,
},
SinglePath {
value: Value,
},
Retained(RetainedSlotRow),
}
impl RowView {
fn missing_required_slot_error(_index: usize) -> InternalError {
InternalError::query_executor_invariant()
}
#[must_use]
#[cfg(test)]
pub(in crate::db::executor) const fn new(slots: Vec<Option<Value>>) -> Self {
Self {
storage: RowViewStorage::Dense(slots),
}
}
#[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 },
}
}
#[cfg(test)]
#[must_use]
pub(in crate::db::executor) fn borrow_slot_for_test(&self, index: usize) -> Option<&Value> {
match &self.storage {
RowViewStorage::Dense(slots) => slots.get(index).and_then(Option::as_ref),
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 slot_value_ref(&self, index: usize) -> Option<&Value> {
match &self.storage {
#[cfg(test)]
RowViewStorage::Dense(slots) => slots.get(index).and_then(Option::as_ref),
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(|| Self::missing_required_slot_error(index))
}
pub(in crate::db::executor) fn into_required_slot_value(
self,
index: usize,
) -> Result<Value, InternalError> {
match self.storage {
#[cfg(test)]
RowViewStorage::Dense(mut slots) => slots
.get_mut(index)
.and_then(Option::take)
.ok_or_else(|| Self::missing_required_slot_error(index)),
RowViewStorage::Single { slot, value } => {
if slot == index {
return Ok(value);
}
Err(Self::missing_required_slot_error(index))
}
RowViewStorage::SinglePath { .. } => Err(Self::missing_required_slot_error(index)),
RowViewStorage::Retained(mut row) => row
.take_slot(index)
.ok_or_else(|| Self::missing_required_slot_error(index)),
}
}
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),
#[cfg(test)]
RowViewStorage::Dense(_) => None,
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),
"grouped row filter expression could not read slot",
)
}
pub(in crate::db::executor) fn extract_orderable_field_value(
&self,
field_slot: FieldSlot,
) -> Result<Value, InternalError> {
let mut value = Some(self.require_slot_owned(field_slot.index)?);
extract_orderable_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,
field: &str,
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, field, segments)?;
Ok(Some(value.map_or(Cow::Owned(Value::Null), Cow::Borrowed)))
}
}
struct SingleGroupedSlotDecode {
slot: usize,
}
struct SingleGroupedPathDecode {
root_slot: usize,
label: String,
segment_bytes: Box<[Box<[u8]>]>,
}
impl SingleGroupedPathDecode {
fn new(path: &ScalarGroupPath) -> Self {
Self {
root_slot: path.root_slot(),
label: path.label().to_string(),
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_grouped_decoded_row(&row, 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_grouped_decoded_row(&row, path.segment_bytes.len())?;
let row_fields = self.row_layout.open_raw_row_with_contract(&row)?;
row_fields.validate_primary_key(key)?;
let root_bytes = row_fields.required_bytes(path.root_slot)?;
let leaf_bytes =
resolve_path_segments(root_bytes, path.segment_bytes.as_ref()).map_err(|_| {
InternalError::persisted_row_field_decode_failed(
path.label.as_str(),
"grouped scalar-path traversal failed",
)
})?;
let value = match leaf_bytes {
Some(leaf_bytes) => {
decode_structural_value_storage_bytes(leaf_bytes).map_err(|_| {
InternalError::persisted_row_field_decode_failed(
path.label.as_str(),
"grouped scalar-path leaf decode failed",
)
})?
}
None => 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,
)?;
let value = value.ok_or_else(InternalError::query_executor_invariant)?;
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<Option<Value>, InternalError> {
charge_grouped_decoded_row(row, 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(crate::db::executor::ExecutorError::missing_row(key).into())
}
}
}
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,
);
}
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()
}
}
fn charge_grouped_decoded_row(row: &RawRow, nested_steps: usize) -> Result<(), InternalError> {
charge_current_execution_budget(
DiagnosticExecutionBudgetResource::DecodedBytes,
u64::try_from(row.len()).unwrap_or(u64::MAX),
)?;
charge_current_execution_budget(
DiagnosticExecutionBudgetResource::NestedValueSteps,
u64::try_from(nested_steps).unwrap_or(u64::MAX),
)
}
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 dense_test_row_view_resolves_sparse_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.borrow_slot_for_test(1), Some(&Value::Nat64(7)));
assert_eq!(
row_view.borrow_slot_for_test(4),
Some(&Value::Text("group".to_string()))
);
assert_eq!(row_view.borrow_slot_for_test(0), None);
}
#[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.borrow_slot_for_test(4),
Some(&Value::Text("group".to_string()))
);
assert_eq!(row_view.borrow_slot_for_test(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",
);
}
}