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