Skip to main content

reifydb_engine/bulk_insert/
builder.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::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		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}