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