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::{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}