1use std::{
5 collections::{HashMap, HashSet},
6 marker::PhantomData,
7};
8
9use reifydb_catalog::{
10 catalog::Catalog,
11 error::{CatalogError, CatalogObjectKind},
12};
13use reifydb_codec::row::{
14 bytes::RowBuilder, pod::EncodedPodRow, series::EncodedSeriesRowBuilder, shape::RowShape,
15 table::EncodedTableRowBuilder,
16};
17use reifydb_core::{
18 error::CoreError,
19 interface::catalog::{
20 id::IndexId,
21 key::PrimaryKey,
22 object::ObjectId,
23 ringbuffer::{RingBuffer, RingBufferMetadata},
24 series::{Series, SeriesMetadata},
25 storage::StorageId,
26 table::Table,
27 },
28 internal_error,
29 key::{
30 any::TaggedKey,
31 catalog::IndexEntryKey,
32 row::{PartitionedRowKey, RowKey},
33 series::{PartitionedSeriesRowKey, SeriesRowKey},
34 },
35};
36use reifydb_runtime::context::clock::Clock;
37use reifydb_transaction::{
38 interceptor::series_row::SeriesRowInterceptor,
39 multi::RangeScope,
40 transaction::{Transaction, command::CommandTransaction},
41};
42use reifydb_value::{
43 fragment::Fragment,
44 params::Params,
45 value::{Value, identity::IdentityId, partition::Partition, row_number::RowNumber},
46};
47
48use super::{
49 BulkInsertResult, RingBufferInsertResult, SeriesInsertResult, TableInsertResult,
50 validation::{
51 reorder_rows_unvalidated, reorder_rows_unvalidated_rb, reorder_rows_unvalidated_series,
52 validate_and_coerce_rows, validate_and_coerce_rows_rb, validate_and_coerce_rows_series,
53 },
54};
55use crate::{
56 Result,
57 bulk_insert::storage::{
58 ringbuffer::{PendingRingBufferInsert, RingBufferInsertBuilder},
59 series::{PendingSeriesInsert, SeriesInsertBuilder},
60 table::{PendingTableInsert, TableInsertBuilder},
61 },
62 engine::StandardEngine,
63 partition::resolve_partition,
64 transaction::operation::{
65 dictionary::DictionaryOperations, ringbuffer::RingBufferOperations, table::TableOperations,
66 },
67 vm::instruction::dml::{
68 primary_key,
69 shape::{get_or_create_ringbuffer_shape, get_or_create_series_shape, get_or_create_table_shape},
70 time::resolve_time,
71 },
72};
73
74pub trait ValidationMode: sealed::Sealed + 'static {
75 const VALIDATED: bool;
76
77 fn run<F, R>(txn: &mut CommandTransaction, total_rows: usize, body: F) -> Result<R>
78 where
79 F: FnOnce(&mut CommandTransaction) -> Result<R>;
80}
81
82pub struct Validated;
83impl ValidationMode for Validated {
84 const VALIDATED: bool = true;
85
86 fn run<F, R>(txn: &mut CommandTransaction, total_rows: usize, body: F) -> Result<R>
87 where
88 F: FnOnce(&mut CommandTransaction) -> Result<R>,
89 {
90 run_checked(txn, total_rows, body)
91 }
92}
93
94pub struct Unchecked;
95impl ValidationMode for Unchecked {
96 const VALIDATED: bool = false;
97
98 fn run<F, R>(txn: &mut CommandTransaction, _total_rows: usize, body: F) -> Result<R>
99 where
100 F: FnOnce(&mut CommandTransaction) -> Result<R>,
101 {
102 txn.execute_bulk_unchecked(body)
103 }
104}
105
106fn run_checked<F, R>(txn: &mut CommandTransaction, total_rows: usize, body: F) -> Result<R>
107where
108 F: FnOnce(&mut CommandTransaction) -> Result<R>,
109{
110 if total_rows > 0 {
111 txn.reserve_writes(total_rows.saturating_mul(2))?;
112 }
113 let r = body(txn)?;
114 txn.commit()?;
115 Ok(r)
116}
117
118pub mod sealed {
119
120 use super::{Unchecked, Validated};
121 pub trait Sealed {}
122 impl Sealed for Validated {}
123 impl Sealed for Unchecked {}
124}
125
126pub struct BulkInsertBuilder<'e, V: ValidationMode = Validated> {
127 engine: &'e StandardEngine,
128 identity: IdentityId,
129 pending_tables: Vec<PendingTableInsert>,
130 pending_ringbuffers: Vec<PendingRingBufferInsert>,
131 pending_series: Vec<PendingSeriesInsert>,
132 _validation: PhantomData<V>,
133}
134
135impl<'e> BulkInsertBuilder<'e, Validated> {
136 pub(crate) fn new(engine: &'e StandardEngine, identity: IdentityId) -> Self {
137 Self {
138 engine,
139 identity,
140 pending_tables: Vec::new(),
141 pending_ringbuffers: Vec::new(),
142 pending_series: Vec::new(),
143 _validation: PhantomData,
144 }
145 }
146}
147
148impl<'e> BulkInsertBuilder<'e, Unchecked> {
149 pub(crate) fn new_unchecked(engine: &'e StandardEngine, identity: IdentityId) -> Self {
150 Self {
151 engine,
152 identity,
153 pending_tables: Vec::new(),
154 pending_ringbuffers: Vec::new(),
155 pending_series: Vec::new(),
156 _validation: PhantomData,
157 }
158 }
159}
160
161impl<'e, V: ValidationMode> BulkInsertBuilder<'e, V> {
162 pub fn table<'a>(&'a mut self, qualified_name: &str) -> TableInsertBuilder<'a, 'e, V> {
163 let (namespace, table) = parse_qualified_name(qualified_name);
164 TableInsertBuilder::new(self, namespace, table)
165 }
166
167 pub fn ringbuffer<'a>(&'a mut self, qualified_name: &str) -> RingBufferInsertBuilder<'a, 'e, V> {
168 let (namespace, ringbuffer) = parse_qualified_name(qualified_name);
169 RingBufferInsertBuilder::new(self, namespace, ringbuffer)
170 }
171
172 pub fn series<'a>(&'a mut self, qualified_name: &str) -> SeriesInsertBuilder<'a, 'e, V> {
173 let (namespace, series) = parse_qualified_name(qualified_name);
174 SeriesInsertBuilder::new(self, namespace, series)
175 }
176
177 pub(super) fn add_table_insert(&mut self, pending: PendingTableInsert) {
178 self.pending_tables.push(pending);
179 }
180
181 pub(super) fn add_ringbuffer_insert(&mut self, pending: PendingRingBufferInsert) {
182 self.pending_ringbuffers.push(pending);
183 }
184
185 pub(super) fn add_series_insert(&mut self, pending: PendingSeriesInsert) {
186 self.pending_series.push(pending);
187 }
188
189 pub fn execute(self) -> Result<BulkInsertResult> {
190 self.engine.reject_if_read_only()?;
191 let mut txn = self.engine.begin_command(self.identity)?;
192 let catalog = self.engine.catalog();
193 let clock = self.engine.clock();
194 let total_rows = self.total_pending_rows();
195 let pending_tables = self.pending_tables;
196 let pending_ringbuffers = self.pending_ringbuffers;
197 let pending_series = self.pending_series;
198
199 V::run(&mut txn, total_rows, move |txn| {
200 run_all_pending::<V>(catalog, clock, txn, pending_tables, pending_ringbuffers, pending_series)
201 })
202 }
203
204 #[inline]
205 fn total_pending_rows(&self) -> usize {
206 self.pending_tables.iter().map(|p| p.rows.len()).sum::<usize>()
207 + self.pending_ringbuffers.iter().map(|p| p.rows.len()).sum::<usize>()
208 + self.pending_series.iter().map(|p| p.rows.len()).sum::<usize>()
209 }
210}
211
212#[inline]
213fn run_all_pending<V: ValidationMode>(
214 catalog: Catalog,
215 clock: &Clock,
216 txn: &mut CommandTransaction,
217 pending_tables: Vec<PendingTableInsert>,
218 pending_ringbuffers: Vec<PendingRingBufferInsert>,
219 pending_series: Vec<PendingSeriesInsert>,
220) -> Result<BulkInsertResult> {
221 let mut result = BulkInsertResult::default();
222 for pending in pending_tables {
223 result.tables.push(execute_table_insert::<V>(&catalog, txn, &pending, clock)?);
224 }
225 for pending in pending_ringbuffers {
226 result.ringbuffers.push(execute_ringbuffer_insert::<V>(&catalog, txn, &pending, clock)?);
227 }
228 for pending in pending_series {
229 result.series.push(execute_series_insert::<V>(&catalog, txn, &pending, clock)?);
230 }
231 Ok(result)
232}
233
234fn execute_table_insert<V: ValidationMode>(
235 catalog: &Catalog,
236 txn: &mut CommandTransaction,
237 pending: &PendingTableInsert,
238 clock: &Clock,
239) -> Result<TableInsertResult> {
240 let table = resolve_table(catalog, txn, pending)?;
241 let shape = get_or_create_table_shape(catalog, &table, &mut Transaction::Command(txn))?;
242 let encoded_bytes_list = encode_table_rows::<V>(catalog, txn, pending, &table, &shape, clock)?;
243 if encoded_bytes_list.is_empty() {
244 return Ok(empty_table_result(pending));
245 }
246 write_table_rows(catalog, txn, &table, &shape, pending, encoded_bytes_list)
247}
248
249#[inline]
250fn empty_table_result(pending: &PendingTableInsert) -> TableInsertResult {
251 TableInsertResult {
252 namespace: pending.namespace.clone(),
253 table: pending.table.clone(),
254 inserted: 0,
255 }
256}
257
258#[inline]
259fn write_table_rows(
260 catalog: &Catalog,
261 txn: &mut CommandTransaction,
262 table: &Table,
263 shape: &RowShape,
264 pending: &PendingTableInsert,
265 encoded_bytes_list: Vec<EncodedTableRowBuilder>,
266) -> Result<TableInsertResult> {
267 let total_rows = encoded_bytes_list.len();
268 let row_numbers = catalog.next_row_number_batch(txn, table.id, total_rows as u64)?;
269 let pk_def = primary_key::get_primary_key(catalog, &mut Transaction::Command(txn), table)?;
270
271 let mut owned_rows = encoded_bytes_list;
272 txn.insert_table(table, shape, &row_numbers, &mut owned_rows)?;
273
274 if let Some(ref pk_def) = pk_def {
275 for (row, &row_number) in owned_rows.iter().zip(row_numbers.iter()) {
276 write_primary_key_index(txn, table, shape, pk_def, row, row_number)?;
277 }
278 }
279
280 Ok(TableInsertResult {
281 namespace: pending.namespace.clone(),
282 table: pending.table.clone(),
283 inserted: total_rows as u64,
284 })
285}
286
287fn resolve_table(catalog: &Catalog, txn: &mut CommandTransaction, pending: &PendingTableInsert) -> Result<Table> {
288 let namespace = catalog
289 .find_namespace_by_name(&mut Transaction::Command(txn), &pending.namespace)?
290 .ok_or_else(|| CatalogError::NotFound {
291 kind: CatalogObjectKind::Namespace,
292 namespace: pending.namespace.to_string(),
293 name: String::new(),
294 fragment: Fragment::None,
295 })?;
296
297 catalog.find_table_by_name(&mut Transaction::Command(txn), namespace.id(), &pending.table)?.ok_or_else(|| {
298 CatalogError::NotFound {
299 kind: CatalogObjectKind::Table,
300 namespace: pending.namespace.to_string(),
301 name: pending.table.to_string(),
302 fragment: Fragment::None,
303 }
304 .into()
305 })
306}
307
308fn encode_table_rows<V: ValidationMode>(
309 catalog: &Catalog,
310 txn: &mut CommandTransaction,
311 pending: &PendingTableInsert,
312 table: &Table,
313 shape: &RowShape,
314 clock: &Clock,
315) -> Result<Vec<EncodedTableRowBuilder>> {
316 let coerced_rows = coerce_table_rows::<V>(&pending.rows, table, txn.identity)?;
317 let mut encoded_bytes_list = Vec::with_capacity(coerced_rows.len());
318 for values in coerced_rows {
319 encoded_bytes_list.push(prepare_table_row::<V>(catalog, txn, table, shape, clock, values)?);
320 }
321 Ok(encoded_bytes_list)
322}
323
324#[inline]
325fn coerce_table_rows<V: ValidationMode>(
326 rows: &[Params],
327 table: &Table,
328 identity: IdentityId,
329) -> Result<Vec<Vec<Value>>> {
330 if V::VALIDATED {
331 validate_and_coerce_rows(rows, table, identity)
332 } else {
333 reorder_rows_unvalidated(rows, table)
334 }
335}
336
337#[inline]
338fn prepare_table_row<V: ValidationMode>(
339 catalog: &Catalog,
340 txn: &mut CommandTransaction,
341 table: &Table,
342 shape: &RowShape,
343 clock: &Clock,
344 mut values: Vec<Value>,
345) -> Result<EncodedTableRowBuilder> {
346 fill_auto_increment_table(catalog, txn, table, &mut values)?;
347 dictionary_encode_table(catalog, txn, table, &mut values)?;
348 if V::VALIDATED {
349 validate_table_constraints(table, &values)?;
350 }
351 encode_row(table, shape, &values, clock)
352}
353
354#[inline]
355fn validate_table_constraints(table: &Table, values: &[Value]) -> Result<()> {
356 for (idx, col) in table.columns.iter().enumerate() {
357 col.constraint.validate(&values[idx])?;
358 }
359 Ok(())
360}
361
362fn fill_auto_increment_table(
363 catalog: &Catalog,
364 txn: &mut CommandTransaction,
365 table: &Table,
366 values: &mut [Value],
367) -> Result<()> {
368 for (idx, col) in table.columns.iter().enumerate() {
369 if col.auto_increment && matches!(values[idx], Value::None { .. }) {
370 values[idx] = catalog.column_sequence_next_value(txn, table.id, col.id)?;
371 }
372 }
373 Ok(())
374}
375
376fn dictionary_encode_table(
377 catalog: &Catalog,
378 txn: &mut CommandTransaction,
379 table: &Table,
380 values: &mut [Value],
381) -> Result<()> {
382 for (idx, col) in table.columns.iter().enumerate() {
383 if let Some(dict_id) = col.dictionary_id {
384 let dictionary =
385 catalog.find_dictionary(&mut Transaction::Command(txn), dict_id)?.ok_or_else(|| {
386 internal_error!("Dictionary {:?} not found for column {}", dict_id, col.name)
387 })?;
388 let entry_id = if matches!(values[idx], Value::None { .. }) {
389 dictionary.id_type.none()
390 } else {
391 txn.insert_into_dictionary(&dictionary, &values[idx])?
392 };
393 values[idx] = entry_id.to_value();
394 }
395 }
396 Ok(())
397}
398
399fn encode_row(table: &Table, shape: &RowShape, values: &[Value], clock: &Clock) -> Result<EncodedTableRowBuilder> {
400 let mut row = shape.allocate_table();
401 for (idx, value) in values.iter().enumerate() {
402 shape.set_value(&mut row, idx, value);
403 }
404 let now = clock.now();
405 row.set_timestamps(now, now);
406 if let Some(time) = resolve_time(&table.name, &table.columns, &table.time, shape, &row, now)? {
407 row.set_time(time);
408 }
409 Ok(row)
410}
411
412fn write_primary_key_index(
413 txn: &mut CommandTransaction,
414 table: &Table,
415 shape: &RowShape,
416 pk_def: &PrimaryKey,
417 row: &[u8],
418 row_number: RowNumber,
419) -> Result<()> {
420 let index_key = primary_key::encode_primary_key(pk_def, row, table, shape)?;
421 let index_entry_key = IndexEntryKey::new(table.id, IndexId::primary(pk_def.id), index_key);
422
423 if txn.contains(&index_entry_key)? {
424 let key_columns = pk_def.columns.iter().map(|c| c.name.clone()).collect();
425 return Err(CoreError::PrimaryKeyViolation {
426 fragment: Fragment::None,
427 table_name: table.name.clone(),
428 key_columns,
429 }
430 .into());
431 }
432
433 txn.set(&index_entry_key, EncodedPodRow::new(&u64::from(row_number).to_be_bytes()).into_bytes())?;
434 Ok(())
435}
436
437fn execute_ringbuffer_insert<V: ValidationMode>(
438 catalog: &Catalog,
439 txn: &mut CommandTransaction,
440 pending: &PendingRingBufferInsert,
441 clock: &Clock,
442) -> Result<RingBufferInsertResult> {
443 let ringbuffer = resolve_ringbuffer(catalog, txn, pending)?;
444 let shape = get_or_create_ringbuffer_shape(catalog, &ringbuffer, &mut Transaction::Command(txn))?;
445 let coerced_rows = coerce_ringbuffer_rows::<V>(pending, &ringbuffer, txn.identity)?;
446 let inserted = insert_ringbuffer_rows::<V>(catalog, txn, &ringbuffer, &shape, coerced_rows, clock)?;
447 Ok(RingBufferInsertResult {
448 namespace: pending.namespace.clone(),
449 ringbuffer: pending.ringbuffer.clone(),
450 inserted,
451 })
452}
453
454#[inline]
455fn resolve_ringbuffer(
456 catalog: &Catalog,
457 txn: &mut CommandTransaction,
458 pending: &PendingRingBufferInsert,
459) -> Result<RingBuffer> {
460 let namespace = catalog
461 .find_namespace_by_name(&mut Transaction::Command(txn), &pending.namespace)?
462 .ok_or_else(|| CatalogError::NotFound {
463 kind: CatalogObjectKind::Namespace,
464 namespace: pending.namespace.to_string(),
465 name: String::new(),
466 fragment: Fragment::None,
467 })?;
468
469 catalog.find_ringbuffer_by_name(&mut Transaction::Command(txn), namespace.id(), &pending.ringbuffer)?
470 .ok_or_else(|| {
471 CatalogError::NotFound {
472 kind: CatalogObjectKind::RingBuffer,
473 namespace: pending.namespace.to_string(),
474 name: pending.ringbuffer.to_string(),
475 fragment: Fragment::None,
476 }
477 .into()
478 })
479}
480
481#[inline]
482fn coerce_ringbuffer_rows<V: ValidationMode>(
483 pending: &PendingRingBufferInsert,
484 ringbuffer: &RingBuffer,
485 identity: IdentityId,
486) -> Result<Vec<Vec<Value>>> {
487 if V::VALIDATED {
488 validate_and_coerce_rows_rb(&pending.rows, ringbuffer, identity)
489 } else {
490 reorder_rows_unvalidated_rb(&pending.rows, ringbuffer)
491 }
492}
493
494fn insert_ringbuffer_rows<V: ValidationMode>(
495 catalog: &Catalog,
496 txn: &mut CommandTransaction,
497 ringbuffer: &RingBuffer,
498 shape: &RowShape,
499 coerced_rows: Vec<Vec<Value>>,
500 clock: &Clock,
501) -> Result<u64> {
502 let partition_col_indices = compute_ringbuffer_partition_col_indices(ringbuffer);
503 let mut cache: HashMap<Vec<Value>, RingBufferMetadata> = HashMap::new();
504 let mut inserted_count = 0u64;
505
506 for mut values in coerced_rows {
507 dict_encode_ringbuffer_row(catalog, txn, ringbuffer, &mut values)?;
508
509 if V::VALIDATED {
510 for (idx, col) in ringbuffer.columns.iter().enumerate() {
511 col.constraint.validate(&values[idx])?;
512 }
513 }
514
515 let partition_key: Vec<Value> = partition_col_indices.iter().map(|&idx| values[idx].clone()).collect();
516 let partition = if partition_col_indices.is_empty() {
517 None
518 } else {
519 Some(Partition::of(&partition_key))
520 };
521
522 let mut row = shape.allocate_ringbuffer();
523 for (idx, value) in values.iter().enumerate() {
524 shape.set_value(&mut row, idx, value);
525 }
526 let now = clock.now();
527 row.set_timestamps(now, now);
528 if let Some(time) =
529 resolve_time(&ringbuffer.name, &ringbuffer.columns, &ringbuffer.time, shape, &row, now)?
530 {
531 row.set_time(time);
532 }
533
534 ensure_ringbuffer_partition_metadata(catalog, txn, ringbuffer, &partition_key, &mut cache)?;
535 let metadata = cache.get_mut(&partition_key).unwrap();
536
537 if metadata.is_full(ringbuffer.capacity) {
538 evict_oldest_for_partition(txn, ringbuffer, partition, metadata)?;
539 }
540
541 let row_number = catalog.next_row_number_for_ringbuffer(txn, ringbuffer.id)?;
542 txn.insert_ringbuffer_at(ringbuffer, shape, partition, row_number, row.freeze_bytes())?;
543
544 if metadata.is_empty() {
545 metadata.head = row_number.0;
546 }
547 metadata.count += 1;
548 metadata.tail = row_number.0 + 1;
549
550 inserted_count += 1;
551 }
552
553 for (partition_key, metadata) in &cache {
554 if metadata.is_empty() {
555 catalog.remove_partition_metadata(&mut Transaction::Command(txn), ringbuffer, partition_key)?;
556 } else {
557 catalog.save_partition_metadata(
558 &mut Transaction::Command(txn),
559 ringbuffer,
560 partition_key,
561 metadata,
562 )?;
563 }
564 }
565
566 Ok(inserted_count)
567}
568
569#[inline]
570fn compute_ringbuffer_partition_col_indices(ringbuffer: &RingBuffer) -> Vec<usize> {
571 ringbuffer
572 .partition_by
573 .iter()
574 .map(|pb_col| ringbuffer.columns.iter().position(|c| c.name == *pb_col).unwrap())
575 .collect()
576}
577
578#[inline]
579fn ensure_ringbuffer_partition_metadata(
580 catalog: &Catalog,
581 txn: &mut CommandTransaction,
582 ringbuffer: &RingBuffer,
583 partition_key: &[Value],
584 cache: &mut HashMap<Vec<Value>, RingBufferMetadata>,
585) -> Result<()> {
586 if !cache.contains_key(partition_key) {
587 let existing =
588 catalog.find_partition_metadata(&mut Transaction::Command(txn), ringbuffer, partition_key)?;
589 let m = existing.unwrap_or_else(RingBufferMetadata::new);
590 cache.insert(partition_key.to_vec(), m);
591 }
592 Ok(())
593}
594
595fn evict_oldest_for_partition(
596 txn: &mut CommandTransaction,
597 ringbuffer: &RingBuffer,
598 partition: Option<Partition>,
599 metadata: &mut RingBufferMetadata,
600) -> Result<()> {
601 if let Some(partition) = partition {
602 let range = PartitionedRowKey::partition_scan_range(ringbuffer.id, partition, None);
603 let oldest = txn.range_rev(range, RangeScope::All, 1)?.next().transpose()?;
604 if let Some(entry) = oldest
605 && let TaggedKey::PartitionedRow(pk) = &entry.key
606 {
607 txn.remove_from_ringbuffer(ringbuffer, Some(partition), pk.row)?;
608 }
609 metadata.count -= 1;
610 return Ok(());
611 }
612
613 let mut evict_pos = metadata.head;
614 loop {
615 let key = RowKey::new(ringbuffer.id, RowNumber(evict_pos));
616 if txn.get(&key)?.is_some() {
617 txn.remove_from_ringbuffer(ringbuffer, None, RowNumber(evict_pos))?;
618 break;
619 }
620 evict_pos += 1;
621 if evict_pos >= metadata.tail {
622 break;
623 }
624 }
625 metadata.head = evict_pos + 1;
626 while metadata.head < metadata.tail {
627 let key = RowKey::new(ringbuffer.id, RowNumber(metadata.head));
628 if txn.get(&key)?.is_some() {
629 break;
630 }
631 metadata.head += 1;
632 }
633 metadata.count -= 1;
634 Ok(())
635}
636
637#[inline]
638fn dict_encode_ringbuffer_row(
639 catalog: &Catalog,
640 txn: &mut CommandTransaction,
641 ringbuffer: &RingBuffer,
642 values: &mut [Value],
643) -> Result<()> {
644 for (idx, col) in ringbuffer.columns.iter().enumerate() {
645 if let Some(dict_id) = col.dictionary_id {
646 let dictionary =
647 catalog.find_dictionary(&mut Transaction::Command(txn), dict_id)?.ok_or_else(|| {
648 internal_error!("Dictionary {:?} not found for column {}", dict_id, col.name)
649 })?;
650 let entry_id = if matches!(values[idx], Value::None { .. }) {
651 dictionary.id_type.none()
652 } else {
653 txn.insert_into_dictionary(&dictionary, &values[idx])?
654 };
655 values[idx] = entry_id.to_value();
656 }
657 }
658 Ok(())
659}
660
661#[inline]
662fn dict_encode_series_row(
663 catalog: &Catalog,
664 txn: &mut CommandTransaction,
665 series: &Series,
666 values: &mut [Value],
667) -> Result<()> {
668 for (idx, col) in series.columns.iter().enumerate() {
669 if let Some(dict_id) = col.dictionary_id {
670 let dictionary =
671 catalog.find_dictionary(&mut Transaction::Command(txn), dict_id)?.ok_or_else(|| {
672 internal_error!("Dictionary {:?} not found for column {}", dict_id, col.name)
673 })?;
674 let entry_id = if matches!(values[idx], Value::None { .. }) {
675 dictionary.id_type.none()
676 } else {
677 txn.insert_into_dictionary(&dictionary, &values[idx])?
678 };
679 values[idx] = entry_id.to_value();
680 }
681 }
682 Ok(())
683}
684
685fn execute_series_insert<V: ValidationMode>(
686 catalog: &Catalog,
687 txn: &mut CommandTransaction,
688 pending: &PendingSeriesInsert,
689 clock: &Clock,
690) -> Result<SeriesInsertResult> {
691 let series = resolve_series(catalog, txn, pending)?;
692 let mut metadata = load_series_metadata(catalog, txn, pending, &series)?;
693 let shape = get_or_create_series_shape(catalog, &series, &mut Transaction::Command(txn))?;
694 let coerced_rows = coerce_series_rows::<V>(pending, &series, txn.identity)?;
695 let inserted = insert_series_rows::<V>(catalog, txn, &series, &shape, coerced_rows, &mut metadata, clock)?;
696 catalog.update_series_metadata_txn(&mut Transaction::Command(txn), series.id, metadata)?;
697 Ok(SeriesInsertResult {
698 namespace: pending.namespace.clone(),
699 series: pending.series.clone(),
700 inserted,
701 })
702}
703
704#[inline]
705fn resolve_series(catalog: &Catalog, txn: &mut CommandTransaction, pending: &PendingSeriesInsert) -> Result<Series> {
706 let namespace = catalog
707 .find_namespace_by_name(&mut Transaction::Command(txn), &pending.namespace)?
708 .ok_or_else(|| CatalogError::NotFound {
709 kind: CatalogObjectKind::Namespace,
710 namespace: pending.namespace.to_string(),
711 name: String::new(),
712 fragment: Fragment::None,
713 })?;
714
715 catalog.find_series_by_name(&mut Transaction::Command(txn), namespace.id(), &pending.series)?.ok_or_else(|| {
716 CatalogError::NotFound {
717 kind: CatalogObjectKind::Series,
718 namespace: pending.namespace.to_string(),
719 name: pending.series.to_string(),
720 fragment: Fragment::None,
721 }
722 .into()
723 })
724}
725
726#[inline]
727fn load_series_metadata(
728 catalog: &Catalog,
729 txn: &mut CommandTransaction,
730 pending: &PendingSeriesInsert,
731 series: &Series,
732) -> Result<SeriesMetadata> {
733 catalog.find_series_metadata(&mut Transaction::Command(txn), series.id)?.ok_or_else(|| {
734 CatalogError::NotFound {
735 kind: CatalogObjectKind::Series,
736 namespace: pending.namespace.to_string(),
737 name: pending.series.to_string(),
738 fragment: Fragment::None,
739 }
740 .into()
741 })
742}
743
744#[inline]
745fn coerce_series_rows<V: ValidationMode>(
746 pending: &PendingSeriesInsert,
747 series: &Series,
748 identity: IdentityId,
749) -> Result<Vec<Vec<Value>>> {
750 if V::VALIDATED {
751 validate_and_coerce_rows_series(&pending.rows, series, identity)
752 } else {
753 reorder_rows_unvalidated_series(&pending.rows, series)
754 }
755}
756
757fn insert_series_rows<V: ValidationMode>(
758 catalog: &Catalog,
759 txn: &mut CommandTransaction,
760 series: &Series,
761 shape: &RowShape,
762 coerced_rows: Vec<Vec<Value>>,
763 metadata: &mut SeriesMetadata,
764 clock: &Clock,
765) -> Result<u64> {
766 let key_col_name = series.key.column();
767 let key_col_idx =
768 series.columns.iter().position(|c| c.name == key_col_name).ok_or_else(|| {
769 internal_error!("series {} key column {} not found", series.name, key_col_name)
770 })?;
771 let partition_col_indices = series_partition_col_indices(series)?;
772 let storage = StorageId::series(series.id);
773 let mut verified: HashSet<Partition> = HashSet::new();
774
775 let mut inserted_count = 0u64;
776 for mut values in coerced_rows {
777 dict_encode_series_row(catalog, txn, series, &mut values)?;
778
779 if V::VALIDATED {
780 for (idx, col) in series.columns.iter().enumerate() {
781 col.constraint.validate(&values[idx])?;
782 }
783 }
784
785 let key_value = series.key_to_u64(values[key_col_idx].clone()).unwrap_or(0);
786
787 metadata.sequence_counter += 1;
788 let sequence = metadata.sequence_counter;
789 let key: TaggedKey = if partition_col_indices.is_empty() {
790 SeriesRowKey {
791 storage,
792 variant_tag: None,
793 key: key_value,
794 sequence,
795 }
796 .into()
797 } else {
798 let partition_values: Vec<Value> =
799 partition_col_indices.iter().map(|&idx| values[idx].clone()).collect();
800 let partition = Partition::of(&partition_values);
801 resolve_partition(
802 &mut Transaction::Command(txn),
803 ObjectId::Series(series.id),
804 partition,
805 &partition_values,
806 &mut verified,
807 )?;
808 PartitionedSeriesRowKey::new(storage, partition, None, key_value, sequence).into()
809 };
810
811 let row = encode_series_row(series, shape, key_value, &values, key_col_idx, clock)?;
812
813 let mut rows_buf = [row];
814 SeriesRowInterceptor::pre_insert(txn, series, &mut rows_buf)?;
815 let [row] = rows_buf;
816 let row = row.freeze_bytes();
817 txn.set(&key, row.clone())?;
818 let rows = [row.clone()];
819 SeriesRowInterceptor::post_insert(txn, series, &rows)?;
820
821 update_series_metadata_for_insert(metadata, key_value);
822 inserted_count += 1;
823 }
824 Ok(inserted_count)
825}
826
827#[inline]
828fn series_partition_col_indices(series: &Series) -> Result<Vec<usize>> {
829 series.partition_by
830 .iter()
831 .map(|name| {
832 series.columns.iter().position(|c| &c.name == name).ok_or_else(|| {
833 internal_error!("series {} partition column {} not found", series.name, name)
834 })
835 })
836 .collect()
837}
838
839#[inline]
840fn encode_series_row(
841 series: &Series,
842 shape: &RowShape,
843 key_value: u64,
844 values: &[Value],
845 key_col_idx: usize,
846 clock: &Clock,
847) -> Result<EncodedSeriesRowBuilder> {
848 let key_value_encoded = series.key_from_u64(key_value);
849 let mut row = shape.allocate_series();
850 shape.set_value(&mut row, 0, &key_value_encoded);
851 let mut shape_idx = 1;
852 for (col_idx, value) in values.iter().enumerate() {
853 if col_idx == key_col_idx {
854 continue;
855 }
856 shape.set_value(&mut row, shape_idx, value);
857 shape_idx += 1;
858 }
859 let now = clock.now();
860 row.set_timestamps(now, now);
861 if let Some(time) = resolve_time(&series.name, &series.columns, &series.time, shape, &row, now)? {
862 row.set_time(time);
863 }
864 Ok(row)
865}
866
867#[inline]
868fn update_series_metadata_for_insert(metadata: &mut SeriesMetadata, key_value: u64) {
869 if metadata.row_count == 0 {
870 metadata.oldest_key = key_value;
871 metadata.newest_key = key_value;
872 } else {
873 if key_value < metadata.oldest_key {
874 metadata.oldest_key = key_value;
875 }
876 if key_value > metadata.newest_key {
877 metadata.newest_key = key_value;
878 }
879 }
880 metadata.row_count += 1;
881}
882
883fn parse_qualified_name(qualified_name: &str) -> (String, String) {
884 if let Some((ns, name)) = qualified_name.rsplit_once("::") {
885 (ns.to_string(), name.to_string())
886 } else {
887 ("default".to_string(), qualified_name.to_string())
888 }
889}
890
891#[cfg(test)]
892mod tests {
893 use super::*;
894
895 #[test]
896 fn parse_qualified_name_simple() {
897 assert_eq!(parse_qualified_name("table"), ("default".to_string(), "table".to_string()));
898 }
899
900 #[test]
901 fn parse_qualified_name_single_namespace() {
902 assert_eq!(parse_qualified_name("ns::table"), ("ns".to_string(), "table".to_string()));
903 }
904
905 #[test]
906 fn parse_qualified_name_nested_namespace() {
907 assert_eq!(parse_qualified_name("a::b::table"), ("a::b".to_string(), "table".to_string()));
908 }
909
910 #[test]
911 fn parse_qualified_name_deeply_nested_namespace() {
912 assert_eq!(parse_qualified_name("a::b::c::table"), ("a::b::c".to_string(), "table".to_string()));
913 }
914
915 #[test]
916 fn parse_qualified_name_empty_string() {
917 assert_eq!(parse_qualified_name(""), ("default".to_string(), "".to_string()));
918 }
919}