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