1use 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 ¶ms,
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}