reifydb-engine 0.9.1

Query execution and processing engine for ReifyDB
Documentation
// SPDX-License-Identifier: Apache-2.0
// Copyright (c) 2026 ReifyDB

use std::{collections::HashSet, sync::Arc};

use reifydb_codec::row::bytes::EncodedBytes;
use reifydb_core::{
	error::diagnostic::{
		catalog::{namespace_not_found, ringbuffer_not_found},
		engine,
	},
	interface::{
		catalog::{
			config::{ConfigKey, GetConfig},
			namespace::Namespace,
			policy::{DataOp, PolicyTargetType},
			ringbuffer::{RingBuffer, RingBufferMetadata},
		},
		resolved::{ResolvedNamespace, ResolvedObject, ResolvedRingBuffer},
	},
	key::{
		any::TaggedKey,
		row::{PartitionedRowKey, RowKey},
	},
	value::column::columns::Columns,
};
use reifydb_evaluate::stack::SymbolTable;
use reifydb_rql::{nodes::DeleteRingBufferNode, query::QueryPlan};
use reifydb_transaction::{multi::RangeScope, transaction::Transaction};
use reifydb_value::{
	fragment::Fragment,
	params::Params,
	return_error,
	value::{Value, partition::Partition, row_number::RowNumber},
};

use super::{
	context::{RingBufferTarget, WriteExecCtx},
	returning::{decode_returning_dictionaries, decode_rows_to_columns, evaluate_returning, with_pre_image},
	shape::get_or_create_ringbuffer_shape,
};
use crate::{
	Result,
	policy::PolicyEvaluator,
	transaction::operation::ringbuffer::{RingBufferOperations, apply_ringbuffer_partition_metadata_after_delete},
	vm::{
		services::Services,
		volcano::{
			compile::compile,
			query::{QueryContext, QueryNode, query_budget},
		},
	},
};

pub(crate) fn delete_ringbuffer(
	services: &Arc<Services>,
	txn: &mut Transaction<'_>,
	plan: DeleteRingBufferNode,
	params: Params,
	symbols: &SymbolTable,
) -> Result<Columns> {
	let DeleteRingBufferNode {
		input,
		target,
		returning,
	} = plan;
	let (namespace, ringbuffer) = resolve_delete_ringbuffer_target(services, txn, &target)?;
	let resolved_source = build_delete_ringbuffer_resolved_source(&namespace, &ringbuffer);
	let target_data = RingBufferTarget {
		namespace: &namespace,
		ringbuffer: &ringbuffer,
	};
	let partition_col_indices = compute_partition_col_indices(&ringbuffer);
	let shape = get_or_create_ringbuffer_shape(&services.catalog, &ringbuffer, txn)?;

	let exec = WriteExecCtx {
		services,
		symbols,
	};
	let row_numbers_filter = if let Some(input_plan) = input {
		Some(collect_row_numbers_for_ringbuffer_delete(
			&exec,
			txn,
			*input_plan,
			&target_data,
			&resolved_source,
			&params,
		)?)
	} else {
		None
	};

	let (deleted_count, returned_rows) = delete_ringbuffer_partitions(
		services,
		txn,
		&target_data,
		&partition_col_indices,
		row_numbers_filter.as_ref(),
		returning.is_some(),
	)?;

	if let Some(returning_exprs) = &returning {
		let mut columns = decode_rows_to_columns(&shape, &returned_rows);
		decode_returning_dictionaries(services, txn, &ringbuffer.columns, &mut columns)?;
		let columns = with_pre_image(columns.clone(), &columns);
		return evaluate_returning(services, symbols, returning_exprs, columns, txn.identity());
	}
	Ok(delete_ringbuffer_result(namespace.name(), &ringbuffer.name, deleted_count))
}

#[inline]
fn resolve_delete_ringbuffer_target(
	services: &Arc<Services>,
	txn: &mut Transaction<'_>,
	target: &ResolvedRingBuffer,
) -> Result<(Namespace, RingBuffer)> {
	let namespace_name = target.namespace().name();
	let Some(namespace) = services.catalog.find_namespace_by_name(txn, namespace_name)? else {
		return_error!(namespace_not_found(Fragment::internal(namespace_name), namespace_name));
	};
	let ringbuffer_name = target.name();
	let Some(ringbuffer) = services.catalog.find_ringbuffer_by_name(txn, namespace.id(), ringbuffer_name)? else {
		let fragment = Fragment::internal(target.name());
		return_error!(ringbuffer_not_found(fragment.clone(), namespace_name, ringbuffer_name));
	};
	Ok((namespace, ringbuffer))
}

#[inline]
fn build_delete_ringbuffer_resolved_source(namespace: &Namespace, ringbuffer: &RingBuffer) -> Option<ResolvedObject> {
	let namespace_ident = Fragment::internal(namespace.name());
	let resolved_namespace = ResolvedNamespace::new(namespace_ident, namespace.clone());
	let rb_ident = Fragment::internal(ringbuffer.name.clone());
	let resolved_rb = ResolvedRingBuffer::new(rb_ident, resolved_namespace, ringbuffer.clone());
	Some(ResolvedObject::RingBuffer(resolved_rb))
}

#[inline]
fn compute_partition_col_indices(ringbuffer: &RingBuffer) -> Vec<usize> {
	ringbuffer
		.partition_by
		.iter()
		.map(|pb_col| ringbuffer.columns.iter().position(|c| c.name == *pb_col).unwrap())
		.collect()
}

fn collect_row_numbers_for_ringbuffer_delete(
	exec: &WriteExecCtx<'_>,
	txn: &mut Transaction<'_>,
	input_plan: QueryPlan,
	target: &RingBufferTarget<'_>,
	resolved_source: &Option<ResolvedObject>,
	params: &Params,
) -> Result<HashSet<RowNumber>> {
	let mut row_numbers_to_delete = HashSet::new();

	let batch_size = exec.services.catalog.get_config_uint2(ConfigKey::QueryRowBatchSize) as u64;
	let mut input_node = compile(
		input_plan,
		txn,
		Arc::new(QueryContext {
			services: exec.services.clone(),
			source: resolved_source.clone(),
			batch_size,
			params: params.clone(),
			symbols: exec.symbols.clone(),
			identity: txn.identity(),
			memory: query_budget(exec.services),
		}),
	);

	let context = QueryContext {
		services: exec.services.clone(),
		source: None,
		batch_size,
		params: params.clone(),
		symbols: exec.symbols.clone(),
		identity: txn.identity(),
		memory: query_budget(exec.services),
	};
	input_node.initialize(txn, &context)?;

	let mut mutable_context = context.clone();
	while let Some(columns) = input_node.next(txn, &mut mutable_context)? {
		PolicyEvaluator::new(exec.services, exec.symbols).enforce_write_policies(
			txn,
			target.namespace.name(),
			&target.ringbuffer.name,
			DataOp::Delete,
			&columns,
			PolicyTargetType::RingBuffer,
		)?;
		if columns.row_numbers().is_empty() {
			return_error!(engine::missing_row_number_column());
		}
		row_numbers_to_delete.extend(columns.row_numbers().iter().copied());
	}
	Ok(row_numbers_to_delete)
}

fn delete_ringbuffer_partitions(
	services: &Arc<Services>,
	txn: &mut Transaction<'_>,
	target: &RingBufferTarget<'_>,
	partition_col_indices: &[usize],
	row_numbers_filter: Option<&HashSet<RowNumber>>,
	has_returning: bool,
) -> Result<(u64, Vec<(RowNumber, EncodedBytes)>)> {
	let ringbuffer = target.ringbuffer;
	let mut deleted_count = 0u64;
	let mut returned_rows: Vec<(RowNumber, EncodedBytes)> = Vec::new();
	let partitions = services.catalog.list_ringbuffer_partitions(txn, ringbuffer)?;

	for partition_info in partitions {
		let partition_key = partition_info.partition_values.clone();
		let partition = partition_info.metadata;
		let mut min_remaining_row: Option<u64> = None;
		let mut partition_deleted = 0u64;
		let partition_hash = if partition_col_indices.is_empty() {
			None
		} else {
			Some(Partition::of(&partition_key))
		};

		for row_num in collect_partition_row_numbers(txn, ringbuffer, partition_hash, &partition)? {
			let should_delete = match row_numbers_filter {
				Some(filter) => filter.contains(&row_num),
				None => true,
			};
			if should_delete {
				let deleted_values = txn.remove_from_ringbuffer(ringbuffer, partition_hash, row_num)?;
				if has_returning {
					returned_rows.push((row_num, deleted_values));
				}
				partition_deleted += 1;
				deleted_count += 1;
			} else {
				min_remaining_row =
					Some(min_remaining_row.map_or(row_num.0, |m: u64| m.min(row_num.0)));
			}
		}

		if row_numbers_filter.is_some() {
			apply_ringbuffer_partition_metadata_after_delete(
				&services.catalog,
				txn,
				ringbuffer,
				&partition_key,
				partition,
				partition_deleted,
				min_remaining_row,
			)?;
		} else {
			services.catalog.remove_partition_metadata(txn, ringbuffer, &partition_key)?;
		}
	}
	Ok((deleted_count, returned_rows))
}

#[inline]
fn delete_ringbuffer_result(namespace: &str, ringbuffer: &str, deleted: u64) -> Columns {
	Columns::single_row([
		("namespace", Value::Utf8(namespace.to_string())),
		("ringbuffer", Value::Utf8(ringbuffer.to_string())),
		("deleted", Value::Uint8(deleted)),
	])
}

fn collect_partition_row_numbers(
	txn: &mut Transaction<'_>,
	ringbuffer: &RingBuffer,
	partition: Option<Partition>,
	metadata: &RingBufferMetadata,
) -> Result<Vec<RowNumber>> {
	let Some(partition) = partition else {
		let mut out = Vec::new();
		for row_num_value in metadata.head..metadata.tail {
			let row_num = RowNumber(row_num_value);
			if txn.get(&RowKey::new(ringbuffer.id, row_num))?.is_some() {
				out.push(row_num);
			}
		}
		return Ok(out);
	};

	let mut out = Vec::new();
	let mut last_key = None;
	loop {
		let batch: Vec<_> = txn
			.range(
				PartitionedRowKey::partition_scan_range(ringbuffer.id, partition, last_key.as_ref()),
				RangeScope::All,
				1024,
			)?
			.collect::<Result<Vec<_>>>()?;
		if batch.is_empty() {
			break;
		}
		let n = batch.len();
		for entry in batch {
			if let TaggedKey::PartitionedRow(pk) = &entry.key {
				out.push(pk.row);
			}
			last_key = Some(entry.key.clone());
		}
		if n < 1024 {
			break;
		}
	}

	out.sort_by_key(|rn| rn.0);
	Ok(out)
}