reifydb_engine/bulk_insert/
validation.rs1use 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}