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