Skip to main content

reifydb_engine/vm/instruction/dml/
ringbuffer_delete.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::{collections::HashSet, sync::Arc};
5
6use reifydb_codec::row::bytes::EncodedBytes;
7use reifydb_core::{
8	error::diagnostic::{
9		catalog::{namespace_not_found, ringbuffer_not_found},
10		engine,
11	},
12	interface::{
13		catalog::{
14			config::{ConfigKey, GetConfig},
15			namespace::Namespace,
16			policy::{DataOp, PolicyTargetType},
17			ringbuffer::{RingBuffer, RingBufferMetadata},
18		},
19		resolved::{ResolvedNamespace, ResolvedObject, ResolvedRingBuffer},
20	},
21	key::{
22		any::TaggedKey,
23		row::{PartitionedRowKey, RowKey},
24	},
25	value::column::columns::Columns,
26};
27use reifydb_evaluate::stack::SymbolTable;
28use reifydb_rql::{nodes::DeleteRingBufferNode, query::QueryPlan};
29use reifydb_transaction::{multi::RangeScope, transaction::Transaction};
30use reifydb_value::{
31	fragment::Fragment,
32	params::Params,
33	return_error,
34	value::{Value, partition::Partition, row_number::RowNumber},
35};
36
37use super::{
38	context::{RingBufferTarget, WriteExecCtx},
39	returning::{decode_returning_dictionaries, decode_rows_to_columns, evaluate_returning, with_pre_image},
40	shape::get_or_create_ringbuffer_shape,
41};
42use crate::{
43	Result,
44	policy::PolicyEvaluator,
45	transaction::operation::ringbuffer::{RingBufferOperations, apply_ringbuffer_partition_metadata_after_delete},
46	vm::{
47		services::Services,
48		volcano::{
49			compile::compile,
50			query::{QueryContext, QueryNode, query_budget},
51		},
52	},
53};
54
55pub(crate) fn delete_ringbuffer(
56	services: &Arc<Services>,
57	txn: &mut Transaction<'_>,
58	plan: DeleteRingBufferNode,
59	params: Params,
60	symbols: &SymbolTable,
61) -> Result<Columns> {
62	let DeleteRingBufferNode {
63		input,
64		target,
65		returning,
66	} = plan;
67	let (namespace, ringbuffer) = resolve_delete_ringbuffer_target(services, txn, &target)?;
68	let resolved_source = build_delete_ringbuffer_resolved_source(&namespace, &ringbuffer);
69	let target_data = RingBufferTarget {
70		namespace: &namespace,
71		ringbuffer: &ringbuffer,
72	};
73	let partition_col_indices = compute_partition_col_indices(&ringbuffer);
74	let shape = get_or_create_ringbuffer_shape(&services.catalog, &ringbuffer, txn)?;
75
76	let exec = WriteExecCtx {
77		services,
78		symbols,
79	};
80	let row_numbers_filter = if let Some(input_plan) = input {
81		Some(collect_row_numbers_for_ringbuffer_delete(
82			&exec,
83			txn,
84			*input_plan,
85			&target_data,
86			&resolved_source,
87			&params,
88		)?)
89	} else {
90		None
91	};
92
93	let (deleted_count, returned_rows) = delete_ringbuffer_partitions(
94		services,
95		txn,
96		&target_data,
97		&partition_col_indices,
98		row_numbers_filter.as_ref(),
99		returning.is_some(),
100	)?;
101
102	if let Some(returning_exprs) = &returning {
103		let mut columns = decode_rows_to_columns(&shape, &returned_rows);
104		decode_returning_dictionaries(services, txn, &ringbuffer.columns, &mut columns)?;
105		let columns = with_pre_image(columns.clone(), &columns);
106		return evaluate_returning(services, symbols, returning_exprs, columns, txn.identity());
107	}
108	Ok(delete_ringbuffer_result(namespace.name(), &ringbuffer.name, deleted_count))
109}
110
111#[inline]
112fn resolve_delete_ringbuffer_target(
113	services: &Arc<Services>,
114	txn: &mut Transaction<'_>,
115	target: &ResolvedRingBuffer,
116) -> Result<(Namespace, RingBuffer)> {
117	let namespace_name = target.namespace().name();
118	let Some(namespace) = services.catalog.find_namespace_by_name(txn, namespace_name)? else {
119		return_error!(namespace_not_found(Fragment::internal(namespace_name), namespace_name));
120	};
121	let ringbuffer_name = target.name();
122	let Some(ringbuffer) = services.catalog.find_ringbuffer_by_name(txn, namespace.id(), ringbuffer_name)? else {
123		let fragment = Fragment::internal(target.name());
124		return_error!(ringbuffer_not_found(fragment.clone(), namespace_name, ringbuffer_name));
125	};
126	Ok((namespace, ringbuffer))
127}
128
129#[inline]
130fn build_delete_ringbuffer_resolved_source(namespace: &Namespace, ringbuffer: &RingBuffer) -> Option<ResolvedObject> {
131	let namespace_ident = Fragment::internal(namespace.name());
132	let resolved_namespace = ResolvedNamespace::new(namespace_ident, namespace.clone());
133	let rb_ident = Fragment::internal(ringbuffer.name.clone());
134	let resolved_rb = ResolvedRingBuffer::new(rb_ident, resolved_namespace, ringbuffer.clone());
135	Some(ResolvedObject::RingBuffer(resolved_rb))
136}
137
138#[inline]
139fn compute_partition_col_indices(ringbuffer: &RingBuffer) -> Vec<usize> {
140	ringbuffer
141		.partition_by
142		.iter()
143		.map(|pb_col| ringbuffer.columns.iter().position(|c| c.name == *pb_col).unwrap())
144		.collect()
145}
146
147fn collect_row_numbers_for_ringbuffer_delete(
148	exec: &WriteExecCtx<'_>,
149	txn: &mut Transaction<'_>,
150	input_plan: QueryPlan,
151	target: &RingBufferTarget<'_>,
152	resolved_source: &Option<ResolvedObject>,
153	params: &Params,
154) -> Result<HashSet<RowNumber>> {
155	let mut row_numbers_to_delete = HashSet::new();
156
157	let batch_size = exec.services.catalog.get_config_uint2(ConfigKey::QueryRowBatchSize) as u64;
158	let mut input_node = compile(
159		input_plan,
160		txn,
161		Arc::new(QueryContext {
162			services: exec.services.clone(),
163			source: resolved_source.clone(),
164			batch_size,
165			params: params.clone(),
166			symbols: exec.symbols.clone(),
167			identity: txn.identity(),
168			memory: query_budget(exec.services),
169		}),
170	);
171
172	let context = QueryContext {
173		services: exec.services.clone(),
174		source: None,
175		batch_size,
176		params: params.clone(),
177		symbols: exec.symbols.clone(),
178		identity: txn.identity(),
179		memory: query_budget(exec.services),
180	};
181	input_node.initialize(txn, &context)?;
182
183	let mut mutable_context = context.clone();
184	while let Some(columns) = input_node.next(txn, &mut mutable_context)? {
185		PolicyEvaluator::new(exec.services, exec.symbols).enforce_write_policies(
186			txn,
187			target.namespace.name(),
188			&target.ringbuffer.name,
189			DataOp::Delete,
190			&columns,
191			PolicyTargetType::RingBuffer,
192		)?;
193		if columns.row_numbers().is_empty() {
194			return_error!(engine::missing_row_number_column());
195		}
196		row_numbers_to_delete.extend(columns.row_numbers().iter().copied());
197	}
198	Ok(row_numbers_to_delete)
199}
200
201fn delete_ringbuffer_partitions(
202	services: &Arc<Services>,
203	txn: &mut Transaction<'_>,
204	target: &RingBufferTarget<'_>,
205	partition_col_indices: &[usize],
206	row_numbers_filter: Option<&HashSet<RowNumber>>,
207	has_returning: bool,
208) -> Result<(u64, Vec<(RowNumber, EncodedBytes)>)> {
209	let ringbuffer = target.ringbuffer;
210	let mut deleted_count = 0u64;
211	let mut returned_rows: Vec<(RowNumber, EncodedBytes)> = Vec::new();
212	let partitions = services.catalog.list_ringbuffer_partitions(txn, ringbuffer)?;
213
214	for partition_info in partitions {
215		let partition_key = partition_info.partition_values.clone();
216		let partition = partition_info.metadata;
217		let mut min_remaining_row: Option<u64> = None;
218		let mut partition_deleted = 0u64;
219		let partition_hash = if partition_col_indices.is_empty() {
220			None
221		} else {
222			Some(Partition::of(&partition_key))
223		};
224
225		for row_num in collect_partition_row_numbers(txn, ringbuffer, partition_hash, &partition)? {
226			let should_delete = match row_numbers_filter {
227				Some(filter) => filter.contains(&row_num),
228				None => true,
229			};
230			if should_delete {
231				let deleted_values = txn.remove_from_ringbuffer(ringbuffer, partition_hash, row_num)?;
232				if has_returning {
233					returned_rows.push((row_num, deleted_values));
234				}
235				partition_deleted += 1;
236				deleted_count += 1;
237			} else {
238				min_remaining_row =
239					Some(min_remaining_row.map_or(row_num.0, |m: u64| m.min(row_num.0)));
240			}
241		}
242
243		if row_numbers_filter.is_some() {
244			apply_ringbuffer_partition_metadata_after_delete(
245				&services.catalog,
246				txn,
247				ringbuffer,
248				&partition_key,
249				partition,
250				partition_deleted,
251				min_remaining_row,
252			)?;
253		} else {
254			services.catalog.remove_partition_metadata(txn, ringbuffer, &partition_key)?;
255		}
256	}
257	Ok((deleted_count, returned_rows))
258}
259
260#[inline]
261fn delete_ringbuffer_result(namespace: &str, ringbuffer: &str, deleted: u64) -> Columns {
262	Columns::single_row([
263		("namespace", Value::Utf8(namespace.to_string())),
264		("ringbuffer", Value::Utf8(ringbuffer.to_string())),
265		("deleted", Value::Uint8(deleted)),
266	])
267}
268
269fn collect_partition_row_numbers(
270	txn: &mut Transaction<'_>,
271	ringbuffer: &RingBuffer,
272	partition: Option<Partition>,
273	metadata: &RingBufferMetadata,
274) -> Result<Vec<RowNumber>> {
275	let Some(partition) = partition else {
276		let mut out = Vec::new();
277		for row_num_value in metadata.head..metadata.tail {
278			let row_num = RowNumber(row_num_value);
279			if txn.get(&RowKey::new(ringbuffer.id, row_num))?.is_some() {
280				out.push(row_num);
281			}
282		}
283		return Ok(out);
284	};
285
286	let mut out = Vec::new();
287	let mut last_key = None;
288	loop {
289		let batch: Vec<_> = txn
290			.range(
291				PartitionedRowKey::partition_scan_range(ringbuffer.id, partition, last_key.as_ref()),
292				RangeScope::All,
293				1024,
294			)?
295			.collect::<Result<Vec<_>>>()?;
296		if batch.is_empty() {
297			break;
298		}
299		let n = batch.len();
300		for entry in batch {
301			if let TaggedKey::PartitionedRow(pk) = &entry.key {
302				out.push(pk.row);
303			}
304			last_key = Some(entry.key.clone());
305		}
306		if n < 1024 {
307			break;
308		}
309	}
310
311	out.sort_by_key(|rn| rn.0);
312	Ok(out)
313}