Skip to main content

reifydb_engine/vm/instruction/dml/
ringbuffer_insert.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::{collections::HashMap, sync::Arc};
5
6use reifydb_codec::row::{
7	bytes::{EncodedBytes, RowBuilder},
8	shape::RowShape,
9};
10use reifydb_core::{
11	error::diagnostic::catalog::{namespace_not_found, ringbuffer_not_found},
12	interface::{
13		catalog::{
14			config::{ConfigKey, GetConfig},
15			namespace::Namespace,
16			policy::{DataOp, PolicyTargetType},
17			ringbuffer::{RingBuffer, RingBufferMetadata},
18		},
19		resolved::{ResolvedColumn, ResolvedNamespace, ResolvedObject, ResolvedRingBuffer},
20	},
21	internal_error,
22	value::column::columns::Columns,
23};
24use reifydb_evaluate::stack::SymbolTable;
25use reifydb_rql::{expression::Expression, nodes::InsertRingBufferNode, query::QueryPlan};
26use reifydb_transaction::transaction::Transaction;
27use reifydb_value::{
28	fragment::Fragment,
29	params::Params,
30	reifydb_assertions, return_error,
31	value::{Value, identity::IdentityId, partition::Partition, row_number::RowNumber},
32};
33use tracing::instrument;
34
35use super::{
36	coerce::coerce_value_to_column_type,
37	context::RingBufferTarget,
38	partition::{
39		compute_partition_col_indices, ensure_partition_metadata, evict_oldest_for_partition,
40		save_all_partition_metadata, update_metadata_after_insert,
41	},
42	returning::{decode_returning_dictionaries, decode_rows_to_columns, evaluate_returning},
43	shape::get_or_create_ringbuffer_shape,
44};
45use crate::{
46	Result,
47	policy::PolicyEvaluator,
48	transaction::operation::{dictionary::DictionaryOperations, ringbuffer::RingBufferOperations},
49	vm::{
50		instruction::dml::time::resolve_time,
51		services::Services,
52		volcano::{
53			compile::compile,
54			query::{QueryContext, QueryNode, query_budget},
55		},
56	},
57};
58
59#[instrument(name = "mutate::ringbuffer::insert", level = "trace", skip_all)]
60pub(crate) fn insert_ringbuffer(
61	services: &Arc<Services>,
62	txn: &mut Transaction<'_>,
63	plan: InsertRingBufferNode,
64	params: Params,
65	symbols: &SymbolTable,
66) -> Result<Columns> {
67	let InsertRingBufferNode {
68		input,
69		target,
70		returning,
71	} = plan;
72	let (namespace, ringbuffer, shape) = resolve_insert_ringbuffer_target_and_shape(services, txn, &target)?;
73	let target_data = RingBufferTarget {
74		namespace: &namespace,
75		ringbuffer: &ringbuffer,
76	};
77	let context = build_insert_ringbuffer_query_context(services, &target_data, &params, symbols, txn.identity());
78	let mut input_node = compile_and_initialize_input(*input, txn, &context)?;
79
80	let mut partition_metadata_cache: HashMap<Vec<Value>, RingBufferMetadata> = HashMap::new();
81	let (inserted_count, returned_rows) = drive_ringbuffer_insert(
82		services,
83		txn,
84		symbols,
85		&target_data,
86		&shape,
87		&context,
88		input_node.as_mut(),
89		returning.is_some(),
90		&mut partition_metadata_cache,
91	)?;
92
93	finalize_ringbuffer_insert(
94		services,
95		txn,
96		&target_data,
97		&shape,
98		symbols,
99		&returning,
100		&partition_metadata_cache,
101		&returned_rows,
102		inserted_count,
103	)
104}
105
106#[inline]
107fn compile_and_initialize_input<'a>(
108	input: QueryPlan,
109	txn: &mut Transaction<'a>,
110	context: &Arc<QueryContext>,
111) -> Result<Box<dyn QueryNode>> {
112	let mut input_node = compile(input, txn, context.clone());
113	input_node.initialize(txn, context)?;
114	Ok(input_node)
115}
116
117#[inline]
118#[allow(clippy::too_many_arguments)]
119fn drive_ringbuffer_insert(
120	services: &Arc<Services>,
121	txn: &mut Transaction<'_>,
122	symbols: &SymbolTable,
123	target_data: &RingBufferTarget<'_>,
124	shape: &RowShape,
125	context: &Arc<QueryContext>,
126	input_node: &mut dyn QueryNode,
127	has_returning: bool,
128	partition_metadata_cache: &mut HashMap<Vec<Value>, RingBufferMetadata>,
129) -> Result<(u64, Vec<(RowNumber, EncodedBytes)>)> {
130	let namespace = target_data.namespace;
131	let ringbuffer = target_data.ringbuffer;
132	let partition_col_indices = compute_partition_col_indices(ringbuffer);
133	let mut inserted_count = 0u64;
134	let mut returned_rows: Vec<(RowNumber, EncodedBytes)> = Vec::new();
135
136	let mut mutable_context = (**context).clone();
137	while let Some(columns) = input_node.next(txn, &mut mutable_context)? {
138		PolicyEvaluator::new(services, symbols).enforce_write_policies(
139			txn,
140			namespace.name(),
141			&ringbuffer.name,
142			DataOp::Insert,
143			&columns,
144			PolicyTargetType::RingBuffer,
145		)?;
146
147		let row_count = columns.row_count();
148		for row_idx in 0..row_count {
149			let (row, row_values) = build_insert_ringbuffer_row(
150				services,
151				txn,
152				target_data,
153				shape,
154				&columns,
155				context,
156				row_idx,
157			)?;
158			let partition_key: Vec<Value> =
159				partition_col_indices.iter().map(|&idx| row_values[idx].clone()).collect();
160			let partition = if partition_col_indices.is_empty() {
161				None
162			} else {
163				Some(Partition::of(&partition_key))
164			};
165			ensure_partition_metadata(
166				services,
167				txn,
168				target_data,
169				&partition_key,
170				partition_metadata_cache,
171			)?;
172			let current_metadata = partition_metadata_cache.get_mut(&partition_key).unwrap();
173
174			if current_metadata.is_full(ringbuffer.capacity) {
175				evict_oldest_for_partition(txn, target_data, partition, current_metadata)?;
176			}
177
178			let row_number = services.catalog.next_row_number_for_ringbuffer(txn, ringbuffer.id)?;
179			let stored_row = txn.insert_ringbuffer_at(ringbuffer, shape, partition, row_number, row)?;
180			if has_returning {
181				returned_rows.push((row_number, stored_row));
182			}
183			update_metadata_after_insert(current_metadata, row_number);
184			inserted_count += 1;
185		}
186	}
187
188	Ok((inserted_count, returned_rows))
189}
190
191#[inline]
192#[allow(clippy::too_many_arguments)]
193fn finalize_ringbuffer_insert(
194	services: &Arc<Services>,
195	txn: &mut Transaction<'_>,
196	target_data: &RingBufferTarget<'_>,
197	shape: &RowShape,
198	symbols: &SymbolTable,
199	returning: &Option<Vec<Expression>>,
200	partition_metadata_cache: &HashMap<Vec<Value>, RingBufferMetadata>,
201	returned_rows: &[(RowNumber, EncodedBytes)],
202	inserted_count: u64,
203) -> Result<Columns> {
204	let ringbuffer = target_data.ringbuffer;
205	save_all_partition_metadata(services, txn, ringbuffer, partition_metadata_cache)?;
206
207	reifydb_assertions! {
208		let returning_rows_match = returning.is_none() || returned_rows.len() as u64 == inserted_count;
209		assert!(
210			returning_rows_match,
211			"ringbuffer insert with a RETURNING clause must capture one stored row per inserted row \
212			 so the returned Columns reflect every insert; captured {} rows but inserted {}",
213			returned_rows.len(),
214			inserted_count
215		);
216	}
217
218	if let Some(returning_exprs) = returning {
219		let mut columns = decode_rows_to_columns(shape, returned_rows);
220		decode_returning_dictionaries(services, txn, &ringbuffer.columns, &mut columns)?;
221		return evaluate_returning(services, symbols, returning_exprs, columns, txn.identity());
222	}
223	Ok(insert_ringbuffer_result(target_data.namespace.name(), &ringbuffer.name, inserted_count))
224}
225
226#[inline]
227fn resolve_insert_ringbuffer_target_and_shape(
228	services: &Arc<Services>,
229	txn: &mut Transaction<'_>,
230	target: &ResolvedRingBuffer,
231) -> Result<(Namespace, RingBuffer, RowShape)> {
232	let namespace_name = target.namespace().name();
233	let Some(namespace) = services.catalog.find_namespace_by_name(txn, namespace_name)? else {
234		return_error!(namespace_not_found(Fragment::internal(namespace_name), namespace_name));
235	};
236	let ringbuffer_name = target.name();
237	let Some(ringbuffer) = services.catalog.find_ringbuffer_by_name(txn, namespace.id(), ringbuffer_name)? else {
238		let fragment = Fragment::internal(target.name());
239		return_error!(ringbuffer_not_found(fragment.clone(), namespace_name, ringbuffer_name));
240	};
241	let shape = get_or_create_ringbuffer_shape(&services.catalog, &ringbuffer, txn)?;
242	Ok((namespace, ringbuffer, shape))
243}
244
245#[inline]
246fn build_insert_ringbuffer_query_context(
247	services: &Arc<Services>,
248	target: &RingBufferTarget<'_>,
249	params: &Params,
250	symbols: &SymbolTable,
251	identity: IdentityId,
252) -> Arc<QueryContext> {
253	let namespace_ident = Fragment::internal(target.namespace.name());
254	let resolved_namespace = ResolvedNamespace::new(namespace_ident, target.namespace.clone());
255	let rb_ident = Fragment::internal(target.ringbuffer.name.clone());
256	let resolved_rb = ResolvedRingBuffer::new(rb_ident, resolved_namespace, target.ringbuffer.clone());
257	Arc::new(QueryContext {
258		services: services.clone(),
259		source: Some(ResolvedObject::RingBuffer(resolved_rb)),
260		batch_size: services.catalog.get_config_uint2(ConfigKey::QueryRowBatchSize) as u64,
261		params: params.clone(),
262		symbols: symbols.clone(),
263		identity,
264		memory: query_budget(services),
265	})
266}
267
268fn build_insert_ringbuffer_row(
269	services: &Arc<Services>,
270	txn: &mut Transaction<'_>,
271	target: &RingBufferTarget<'_>,
272	shape: &RowShape,
273	columns: &Columns,
274	context: &Arc<QueryContext>,
275	row_idx: usize,
276) -> Result<(EncodedBytes, Vec<Value>)> {
277	let mut row = shape.allocate_ringbuffer();
278	let mut row_values: Vec<Value> = Vec::with_capacity(target.ringbuffer.columns.len());
279
280	for (rb_idx, rb_column) in target.ringbuffer.columns.iter().enumerate() {
281		let mut value = if let Some(input_column) = columns.iter().find(|col| col.name() == rb_column.name) {
282			input_column.data().get_value(row_idx)
283		} else {
284			Value::none()
285		};
286
287		let column_ident = columns
288			.iter()
289			.find(|col| col.name() == rb_column.name)
290			.map(|col| col.name().clone())
291			.unwrap_or_else(|| Fragment::internal(&rb_column.name));
292		let resolved_column =
293			ResolvedColumn::new(column_ident.clone(), context.source.clone().unwrap(), rb_column.clone());
294
295		value = coerce_value_to_column_type(value, rb_column.constraint.get_type(), resolved_column, context)?;
296		if let Err(mut e) = rb_column.constraint.validate(&value) {
297			e.0.fragment = column_ident.clone();
298			return Err(e);
299		}
300
301		let value = if let Some(dict_id) = rb_column.dictionary_id {
302			let dictionary = services.catalog.find_dictionary(txn, dict_id)?.ok_or_else(|| {
303				internal_error!("Dictionary {:?} not found for column {}", dict_id, rb_column.name)
304			})?;
305			let entry_id = if matches!(value, Value::None { .. }) {
306				dictionary.id_type.none()
307			} else {
308				txn.insert_into_dictionary(&dictionary, &value)?
309			};
310			entry_id.to_value()
311		} else {
312			value
313		};
314
315		row_values.push(value.clone());
316		shape.set_value(&mut row, rb_idx, &value);
317	}
318
319	let now = services.runtime_context.clock.now();
320	row.set_timestamps(now, now);
321	if let Some(time) = resolve_time(
322		&target.ringbuffer.name,
323		&target.ringbuffer.columns,
324		&target.ringbuffer.time,
325		shape,
326		&row,
327		now,
328	)? {
329		row.set_time(time);
330	}
331	Ok((row.freeze_bytes(), row_values))
332}
333
334#[inline]
335fn insert_ringbuffer_result(namespace: &str, ringbuffer: &str, inserted: u64) -> Columns {
336	Columns::single_row([
337		("namespace", Value::Utf8(namespace.to_string())),
338		("ringbuffer", Value::Utf8(ringbuffer.to_string())),
339		("inserted", Value::Uint8(inserted)),
340	])
341}