Skip to main content

reifydb_engine/bulk_insert/
validation.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::iter;
5
6use reifydb_core::{
7	interface::catalog::{column::Column, ringbuffer::RingBuffer, series::Series, table::Table},
8	value::column::buffer::ColumnBuffer,
9};
10use reifydb_value::{
11	fragment::Fragment,
12	params::Params,
13	value::{Value, identity::IdentityId},
14};
15
16use super::coerce::coerce_columns;
17use crate::{Result, error::EngineError};
18
19pub fn validate_and_coerce_rows(rows: &[Params], table: &Table, identity: IdentityId) -> Result<Vec<Vec<Value>>> {
20	if rows.is_empty() {
21		return Ok(Vec::new());
22	}
23
24	let num_cols = table.columns.len();
25	let num_rows = rows.len();
26
27	let column_data = collect_rows_to_columns(rows, &table.columns, &table.name)?;
28	let coerced_columns = coerce_columns(&column_data, &table.columns, num_rows, identity)?;
29
30	Ok(columns_to_rows(&coerced_columns, num_rows, num_cols))
31}
32
33pub fn validate_and_coerce_rows_rb(
34	rows: &[Params],
35	ringbuffer: &RingBuffer,
36	identity: IdentityId,
37) -> Result<Vec<Vec<Value>>> {
38	if rows.is_empty() {
39		return Ok(Vec::new());
40	}
41
42	let num_cols = ringbuffer.columns.len();
43	let num_rows = rows.len();
44
45	let column_data = collect_rows_to_columns(rows, &ringbuffer.columns, &ringbuffer.name)?;
46	let coerced_columns = coerce_columns(&column_data, &ringbuffer.columns, num_rows, identity)?;
47
48	Ok(columns_to_rows(&coerced_columns, num_rows, num_cols))
49}
50
51pub fn reorder_rows_unvalidated(rows: &[Params], table: &Table) -> Result<Vec<Vec<Value>>> {
52	if rows.is_empty() {
53		return Ok(Vec::new());
54	}
55
56	let num_cols = table.columns.len();
57	let num_rows = rows.len();
58
59	let column_data = collect_rows_to_columns(rows, &table.columns, &table.name)?;
60
61	Ok(columns_to_rows(&column_data, num_rows, num_cols))
62}
63
64pub fn reorder_rows_unvalidated_rb(rows: &[Params], ringbuffer: &RingBuffer) -> Result<Vec<Vec<Value>>> {
65	if rows.is_empty() {
66		return Ok(Vec::new());
67	}
68
69	let num_cols = ringbuffer.columns.len();
70	let num_rows = rows.len();
71
72	let column_data = collect_rows_to_columns(rows, &ringbuffer.columns, &ringbuffer.name)?;
73
74	Ok(columns_to_rows(&column_data, num_rows, num_cols))
75}
76
77pub fn validate_and_coerce_rows_series(
78	rows: &[Params],
79	series: &Series,
80	identity: IdentityId,
81) -> Result<Vec<Vec<Value>>> {
82	if rows.is_empty() {
83		return Ok(Vec::new());
84	}
85
86	let num_cols = series.columns.len();
87	let num_rows = rows.len();
88
89	let column_data = collect_rows_to_columns(rows, &series.columns, &series.name)?;
90	let coerced_columns = coerce_columns(&column_data, &series.columns, num_rows, identity)?;
91
92	Ok(columns_to_rows(&coerced_columns, num_rows, num_cols))
93}
94
95pub fn reorder_rows_unvalidated_series(rows: &[Params], series: &Series) -> Result<Vec<Vec<Value>>> {
96	if rows.is_empty() {
97		return Ok(Vec::new());
98	}
99
100	let num_cols = series.columns.len();
101	let num_rows = rows.len();
102
103	let column_data = collect_rows_to_columns(rows, &series.columns, &series.name)?;
104
105	Ok(columns_to_rows(&column_data, num_rows, num_cols))
106}
107
108fn collect_rows_to_columns(rows: &[Params], columns: &[Column], source_name: &str) -> Result<Vec<ColumnBuffer>> {
109	let num_cols = columns.len();
110	let mut column_data: Vec<ColumnBuffer> =
111		columns.iter().map(|col| ColumnBuffer::none_typed(col.constraint.get_type(), 0)).collect();
112
113	for params in rows {
114		match params {
115			Params::Named(map) => {
116				for (col_idx, col) in columns.iter().enumerate() {
117					let value = map.get(&col.name).cloned().unwrap_or(Value::none());
118					column_data[col_idx].push_value(value);
119				}
120			}
121			Params::Positional(vals) => {
122				if vals.len() > num_cols {
123					return Err(EngineError::BulkInsertTooManyValues {
124						fragment: Fragment::None,
125						expected: num_cols,
126						actual: vals.len(),
127					}
128					.into());
129				}
130				for (col_data, val) in
131					column_data.iter_mut().zip(vals.iter().map(Some).chain(iter::repeat(None)))
132				{
133					col_data.push_value(val.cloned().unwrap_or(Value::none()));
134				}
135			}
136			Params::None => {
137				for col_data in column_data.iter_mut() {
138					col_data.push_none();
139				}
140			}
141		}
142	}
143
144	for params in rows {
145		if let Params::Named(map) = params {
146			for name in map.keys() {
147				if !columns.iter().any(|c| &c.name == name) {
148					return Err(EngineError::BulkInsertColumnNotFound {
149						fragment: Fragment::None,
150						table_name: source_name.to_string(),
151						column: name.to_string(),
152					}
153					.into());
154				}
155			}
156		}
157	}
158
159	Ok(column_data)
160}
161
162fn columns_to_rows(columns: &[ColumnBuffer], num_rows: usize, num_cols: usize) -> Vec<Vec<Value>> {
163	let mut result = Vec::with_capacity(num_rows);
164
165	for row_idx in 0..num_rows {
166		let mut row_values = Vec::with_capacity(num_cols);
167		for col in columns.iter().take(num_cols) {
168			row_values.push(col.get_value(row_idx));
169		}
170		result.push(row_values);
171	}
172
173	result
174}