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