use std::{
collections::{HashMap, HashSet},
marker::PhantomData,
};
use reifydb_catalog::{
catalog::Catalog,
error::{CatalogError, CatalogObjectKind},
};
use reifydb_codec::row::{
bytes::RowBuilder, pod::EncodedPodRow, series::EncodedSeriesRowBuilder, shape::RowShape,
table::EncodedTableRowBuilder,
};
use reifydb_core::{
error::CoreError,
interface::catalog::{
id::IndexId,
key::PrimaryKey,
object::ObjectId,
ringbuffer::{RingBuffer, RingBufferMetadata},
series::{Series, SeriesMetadata},
storage::StorageId,
table::Table,
},
internal_error,
key::{
any::TaggedKey,
catalog::IndexEntryKey,
row::{PartitionedRowKey, RowKey},
series::{PartitionedSeriesRowKey, SeriesRowKey},
},
};
use reifydb_runtime::context::clock::Clock;
use reifydb_transaction::{
interceptor::series_row::SeriesRowInterceptor,
multi::RangeScope,
transaction::{Transaction, command::CommandTransaction},
};
use reifydb_value::{
fragment::Fragment,
params::Params,
value::{Value, identity::IdentityId, partition::Partition, row_number::RowNumber},
};
use super::{
BulkInsertResult, RingBufferInsertResult, SeriesInsertResult, TableInsertResult,
validation::{
reorder_rows_unvalidated, reorder_rows_unvalidated_rb, reorder_rows_unvalidated_series,
validate_and_coerce_rows, validate_and_coerce_rows_rb, validate_and_coerce_rows_series,
},
};
use crate::{
Result,
bulk_insert::storage::{
ringbuffer::{PendingRingBufferInsert, RingBufferInsertBuilder},
series::{PendingSeriesInsert, SeriesInsertBuilder},
table::{PendingTableInsert, TableInsertBuilder},
},
engine::StandardEngine,
partition::resolve_partition,
transaction::operation::{
dictionary::DictionaryOperations, ringbuffer::RingBufferOperations, table::TableOperations,
},
vm::instruction::dml::{
primary_key,
shape::{get_or_create_ringbuffer_shape, get_or_create_series_shape, get_or_create_table_shape},
time::resolve_time,
},
};
pub trait ValidationMode: sealed::Sealed + 'static {
const VALIDATED: bool;
fn run<F, R>(txn: &mut CommandTransaction, total_rows: usize, body: F) -> Result<R>
where
F: FnOnce(&mut CommandTransaction) -> Result<R>;
}
pub struct Validated;
impl ValidationMode for Validated {
const VALIDATED: bool = true;
fn run<F, R>(txn: &mut CommandTransaction, total_rows: usize, body: F) -> Result<R>
where
F: FnOnce(&mut CommandTransaction) -> Result<R>,
{
run_checked(txn, total_rows, body)
}
}
pub struct Unchecked;
impl ValidationMode for Unchecked {
const VALIDATED: bool = false;
fn run<F, R>(txn: &mut CommandTransaction, _total_rows: usize, body: F) -> Result<R>
where
F: FnOnce(&mut CommandTransaction) -> Result<R>,
{
txn.execute_bulk_unchecked(body)
}
}
fn run_checked<F, R>(txn: &mut CommandTransaction, total_rows: usize, body: F) -> Result<R>
where
F: FnOnce(&mut CommandTransaction) -> Result<R>,
{
if total_rows > 0 {
txn.reserve_writes(total_rows.saturating_mul(2))?;
}
let r = body(txn)?;
txn.commit()?;
Ok(r)
}
pub mod sealed {
use super::{Unchecked, Validated};
pub trait Sealed {}
impl Sealed for Validated {}
impl Sealed for Unchecked {}
}
pub struct BulkInsertBuilder<'e, V: ValidationMode = Validated> {
engine: &'e StandardEngine,
identity: IdentityId,
pending_tables: Vec<PendingTableInsert>,
pending_ringbuffers: Vec<PendingRingBufferInsert>,
pending_series: Vec<PendingSeriesInsert>,
_validation: PhantomData<V>,
}
impl<'e> BulkInsertBuilder<'e, Validated> {
pub(crate) fn new(engine: &'e StandardEngine, identity: IdentityId) -> Self {
Self {
engine,
identity,
pending_tables: Vec::new(),
pending_ringbuffers: Vec::new(),
pending_series: Vec::new(),
_validation: PhantomData,
}
}
}
impl<'e> BulkInsertBuilder<'e, Unchecked> {
pub(crate) fn new_unchecked(engine: &'e StandardEngine, identity: IdentityId) -> Self {
Self {
engine,
identity,
pending_tables: Vec::new(),
pending_ringbuffers: Vec::new(),
pending_series: Vec::new(),
_validation: PhantomData,
}
}
}
impl<'e, V: ValidationMode> BulkInsertBuilder<'e, V> {
pub fn table<'a>(&'a mut self, qualified_name: &str) -> TableInsertBuilder<'a, 'e, V> {
let (namespace, table) = parse_qualified_name(qualified_name);
TableInsertBuilder::new(self, namespace, table)
}
pub fn ringbuffer<'a>(&'a mut self, qualified_name: &str) -> RingBufferInsertBuilder<'a, 'e, V> {
let (namespace, ringbuffer) = parse_qualified_name(qualified_name);
RingBufferInsertBuilder::new(self, namespace, ringbuffer)
}
pub fn series<'a>(&'a mut self, qualified_name: &str) -> SeriesInsertBuilder<'a, 'e, V> {
let (namespace, series) = parse_qualified_name(qualified_name);
SeriesInsertBuilder::new(self, namespace, series)
}
pub(super) fn add_table_insert(&mut self, pending: PendingTableInsert) {
self.pending_tables.push(pending);
}
pub(super) fn add_ringbuffer_insert(&mut self, pending: PendingRingBufferInsert) {
self.pending_ringbuffers.push(pending);
}
pub(super) fn add_series_insert(&mut self, pending: PendingSeriesInsert) {
self.pending_series.push(pending);
}
pub fn execute(self) -> Result<BulkInsertResult> {
self.engine.reject_if_read_only()?;
let mut txn = self.engine.begin_command(self.identity)?;
let catalog = self.engine.catalog();
let clock = self.engine.clock();
let total_rows = self.total_pending_rows();
let pending_tables = self.pending_tables;
let pending_ringbuffers = self.pending_ringbuffers;
let pending_series = self.pending_series;
V::run(&mut txn, total_rows, move |txn| {
run_all_pending::<V>(catalog, clock, txn, pending_tables, pending_ringbuffers, pending_series)
})
}
#[inline]
fn total_pending_rows(&self) -> usize {
self.pending_tables.iter().map(|p| p.rows.len()).sum::<usize>()
+ self.pending_ringbuffers.iter().map(|p| p.rows.len()).sum::<usize>()
+ self.pending_series.iter().map(|p| p.rows.len()).sum::<usize>()
}
}
#[inline]
fn run_all_pending<V: ValidationMode>(
catalog: Catalog,
clock: &Clock,
txn: &mut CommandTransaction,
pending_tables: Vec<PendingTableInsert>,
pending_ringbuffers: Vec<PendingRingBufferInsert>,
pending_series: Vec<PendingSeriesInsert>,
) -> Result<BulkInsertResult> {
let mut result = BulkInsertResult::default();
for pending in pending_tables {
result.tables.push(execute_table_insert::<V>(&catalog, txn, &pending, clock)?);
}
for pending in pending_ringbuffers {
result.ringbuffers.push(execute_ringbuffer_insert::<V>(&catalog, txn, &pending, clock)?);
}
for pending in pending_series {
result.series.push(execute_series_insert::<V>(&catalog, txn, &pending, clock)?);
}
Ok(result)
}
fn execute_table_insert<V: ValidationMode>(
catalog: &Catalog,
txn: &mut CommandTransaction,
pending: &PendingTableInsert,
clock: &Clock,
) -> Result<TableInsertResult> {
let table = resolve_table(catalog, txn, pending)?;
let shape = get_or_create_table_shape(catalog, &table, &mut Transaction::Command(txn))?;
let encoded_bytes_list = encode_table_rows::<V>(catalog, txn, pending, &table, &shape, clock)?;
if encoded_bytes_list.is_empty() {
return Ok(empty_table_result(pending));
}
write_table_rows(catalog, txn, &table, &shape, pending, encoded_bytes_list)
}
#[inline]
fn empty_table_result(pending: &PendingTableInsert) -> TableInsertResult {
TableInsertResult {
namespace: pending.namespace.clone(),
table: pending.table.clone(),
inserted: 0,
}
}
#[inline]
fn write_table_rows(
catalog: &Catalog,
txn: &mut CommandTransaction,
table: &Table,
shape: &RowShape,
pending: &PendingTableInsert,
encoded_bytes_list: Vec<EncodedTableRowBuilder>,
) -> Result<TableInsertResult> {
let total_rows = encoded_bytes_list.len();
let row_numbers = catalog.next_row_number_batch(txn, table.id, total_rows as u64)?;
let pk_def = primary_key::get_primary_key(catalog, &mut Transaction::Command(txn), table)?;
let mut owned_rows = encoded_bytes_list;
txn.insert_table(table, shape, &row_numbers, &mut owned_rows)?;
if let Some(ref pk_def) = pk_def {
for (row, &row_number) in owned_rows.iter().zip(row_numbers.iter()) {
write_primary_key_index(txn, table, shape, pk_def, row, row_number)?;
}
}
Ok(TableInsertResult {
namespace: pending.namespace.clone(),
table: pending.table.clone(),
inserted: total_rows as u64,
})
}
fn resolve_table(catalog: &Catalog, txn: &mut CommandTransaction, pending: &PendingTableInsert) -> Result<Table> {
let namespace = catalog
.find_namespace_by_name(&mut Transaction::Command(txn), &pending.namespace)?
.ok_or_else(|| CatalogError::NotFound {
kind: CatalogObjectKind::Namespace,
namespace: pending.namespace.to_string(),
name: String::new(),
fragment: Fragment::None,
})?;
catalog.find_table_by_name(&mut Transaction::Command(txn), namespace.id(), &pending.table)?.ok_or_else(|| {
CatalogError::NotFound {
kind: CatalogObjectKind::Table,
namespace: pending.namespace.to_string(),
name: pending.table.to_string(),
fragment: Fragment::None,
}
.into()
})
}
fn encode_table_rows<V: ValidationMode>(
catalog: &Catalog,
txn: &mut CommandTransaction,
pending: &PendingTableInsert,
table: &Table,
shape: &RowShape,
clock: &Clock,
) -> Result<Vec<EncodedTableRowBuilder>> {
let coerced_rows = coerce_table_rows::<V>(&pending.rows, table, txn.identity)?;
let mut encoded_bytes_list = Vec::with_capacity(coerced_rows.len());
for values in coerced_rows {
encoded_bytes_list.push(prepare_table_row::<V>(catalog, txn, table, shape, clock, values)?);
}
Ok(encoded_bytes_list)
}
#[inline]
fn coerce_table_rows<V: ValidationMode>(
rows: &[Params],
table: &Table,
identity: IdentityId,
) -> Result<Vec<Vec<Value>>> {
if V::VALIDATED {
validate_and_coerce_rows(rows, table, identity)
} else {
reorder_rows_unvalidated(rows, table)
}
}
#[inline]
fn prepare_table_row<V: ValidationMode>(
catalog: &Catalog,
txn: &mut CommandTransaction,
table: &Table,
shape: &RowShape,
clock: &Clock,
mut values: Vec<Value>,
) -> Result<EncodedTableRowBuilder> {
fill_auto_increment_table(catalog, txn, table, &mut values)?;
dictionary_encode_table(catalog, txn, table, &mut values)?;
if V::VALIDATED {
validate_table_constraints(table, &values)?;
}
encode_row(table, shape, &values, clock)
}
#[inline]
fn validate_table_constraints(table: &Table, values: &[Value]) -> Result<()> {
for (idx, col) in table.columns.iter().enumerate() {
col.constraint.validate(&values[idx])?;
}
Ok(())
}
fn fill_auto_increment_table(
catalog: &Catalog,
txn: &mut CommandTransaction,
table: &Table,
values: &mut [Value],
) -> Result<()> {
for (idx, col) in table.columns.iter().enumerate() {
if col.auto_increment && matches!(values[idx], Value::None { .. }) {
values[idx] = catalog.column_sequence_next_value(txn, table.id, col.id)?;
}
}
Ok(())
}
fn dictionary_encode_table(
catalog: &Catalog,
txn: &mut CommandTransaction,
table: &Table,
values: &mut [Value],
) -> Result<()> {
for (idx, col) in table.columns.iter().enumerate() {
if let Some(dict_id) = col.dictionary_id {
let dictionary =
catalog.find_dictionary(&mut Transaction::Command(txn), dict_id)?.ok_or_else(|| {
internal_error!("Dictionary {:?} not found for column {}", dict_id, col.name)
})?;
let entry_id = if matches!(values[idx], Value::None { .. }) {
dictionary.id_type.none()
} else {
txn.insert_into_dictionary(&dictionary, &values[idx])?
};
values[idx] = entry_id.to_value();
}
}
Ok(())
}
fn encode_row(table: &Table, shape: &RowShape, values: &[Value], clock: &Clock) -> Result<EncodedTableRowBuilder> {
let mut row = shape.allocate_table();
for (idx, value) in values.iter().enumerate() {
shape.set_value(&mut row, idx, value);
}
let now = clock.now();
row.set_timestamps(now, now);
if let Some(time) = resolve_time(&table.name, &table.columns, &table.time, shape, &row, now)? {
row.set_time(time);
}
Ok(row)
}
fn write_primary_key_index(
txn: &mut CommandTransaction,
table: &Table,
shape: &RowShape,
pk_def: &PrimaryKey,
row: &[u8],
row_number: RowNumber,
) -> Result<()> {
let index_key = primary_key::encode_primary_key(pk_def, row, table, shape)?;
let index_entry_key = IndexEntryKey::new(table.id, IndexId::primary(pk_def.id), index_key);
if txn.contains(&index_entry_key)? {
let key_columns = pk_def.columns.iter().map(|c| c.name.clone()).collect();
return Err(CoreError::PrimaryKeyViolation {
fragment: Fragment::None,
table_name: table.name.clone(),
key_columns,
}
.into());
}
txn.set(&index_entry_key, EncodedPodRow::new(&u64::from(row_number).to_be_bytes()).into_bytes())?;
Ok(())
}
fn execute_ringbuffer_insert<V: ValidationMode>(
catalog: &Catalog,
txn: &mut CommandTransaction,
pending: &PendingRingBufferInsert,
clock: &Clock,
) -> Result<RingBufferInsertResult> {
let ringbuffer = resolve_ringbuffer(catalog, txn, pending)?;
let shape = get_or_create_ringbuffer_shape(catalog, &ringbuffer, &mut Transaction::Command(txn))?;
let coerced_rows = coerce_ringbuffer_rows::<V>(pending, &ringbuffer, txn.identity)?;
let inserted = insert_ringbuffer_rows::<V>(catalog, txn, &ringbuffer, &shape, coerced_rows, clock)?;
Ok(RingBufferInsertResult {
namespace: pending.namespace.clone(),
ringbuffer: pending.ringbuffer.clone(),
inserted,
})
}
#[inline]
fn resolve_ringbuffer(
catalog: &Catalog,
txn: &mut CommandTransaction,
pending: &PendingRingBufferInsert,
) -> Result<RingBuffer> {
let namespace = catalog
.find_namespace_by_name(&mut Transaction::Command(txn), &pending.namespace)?
.ok_or_else(|| CatalogError::NotFound {
kind: CatalogObjectKind::Namespace,
namespace: pending.namespace.to_string(),
name: String::new(),
fragment: Fragment::None,
})?;
catalog.find_ringbuffer_by_name(&mut Transaction::Command(txn), namespace.id(), &pending.ringbuffer)?
.ok_or_else(|| {
CatalogError::NotFound {
kind: CatalogObjectKind::RingBuffer,
namespace: pending.namespace.to_string(),
name: pending.ringbuffer.to_string(),
fragment: Fragment::None,
}
.into()
})
}
#[inline]
fn coerce_ringbuffer_rows<V: ValidationMode>(
pending: &PendingRingBufferInsert,
ringbuffer: &RingBuffer,
identity: IdentityId,
) -> Result<Vec<Vec<Value>>> {
if V::VALIDATED {
validate_and_coerce_rows_rb(&pending.rows, ringbuffer, identity)
} else {
reorder_rows_unvalidated_rb(&pending.rows, ringbuffer)
}
}
fn insert_ringbuffer_rows<V: ValidationMode>(
catalog: &Catalog,
txn: &mut CommandTransaction,
ringbuffer: &RingBuffer,
shape: &RowShape,
coerced_rows: Vec<Vec<Value>>,
clock: &Clock,
) -> Result<u64> {
let partition_col_indices = compute_ringbuffer_partition_col_indices(ringbuffer);
let mut cache: HashMap<Vec<Value>, RingBufferMetadata> = HashMap::new();
let mut inserted_count = 0u64;
for mut values in coerced_rows {
dict_encode_ringbuffer_row(catalog, txn, ringbuffer, &mut values)?;
if V::VALIDATED {
for (idx, col) in ringbuffer.columns.iter().enumerate() {
col.constraint.validate(&values[idx])?;
}
}
let partition_key: Vec<Value> = partition_col_indices.iter().map(|&idx| values[idx].clone()).collect();
let partition = if partition_col_indices.is_empty() {
None
} else {
Some(Partition::of(&partition_key))
};
let mut row = shape.allocate_ringbuffer();
for (idx, value) in values.iter().enumerate() {
shape.set_value(&mut row, idx, value);
}
let now = clock.now();
row.set_timestamps(now, now);
if let Some(time) =
resolve_time(&ringbuffer.name, &ringbuffer.columns, &ringbuffer.time, shape, &row, now)?
{
row.set_time(time);
}
ensure_ringbuffer_partition_metadata(catalog, txn, ringbuffer, &partition_key, &mut cache)?;
let metadata = cache.get_mut(&partition_key).unwrap();
if metadata.is_full(ringbuffer.capacity) {
evict_oldest_for_partition(txn, ringbuffer, partition, metadata)?;
}
let row_number = catalog.next_row_number_for_ringbuffer(txn, ringbuffer.id)?;
txn.insert_ringbuffer_at(ringbuffer, shape, partition, row_number, row.freeze_bytes())?;
if metadata.is_empty() {
metadata.head = row_number.0;
}
metadata.count += 1;
metadata.tail = row_number.0 + 1;
inserted_count += 1;
}
for (partition_key, metadata) in &cache {
if metadata.is_empty() {
catalog.remove_partition_metadata(&mut Transaction::Command(txn), ringbuffer, partition_key)?;
} else {
catalog.save_partition_metadata(
&mut Transaction::Command(txn),
ringbuffer,
partition_key,
metadata,
)?;
}
}
Ok(inserted_count)
}
#[inline]
fn compute_ringbuffer_partition_col_indices(ringbuffer: &RingBuffer) -> Vec<usize> {
ringbuffer
.partition_by
.iter()
.map(|pb_col| ringbuffer.columns.iter().position(|c| c.name == *pb_col).unwrap())
.collect()
}
#[inline]
fn ensure_ringbuffer_partition_metadata(
catalog: &Catalog,
txn: &mut CommandTransaction,
ringbuffer: &RingBuffer,
partition_key: &[Value],
cache: &mut HashMap<Vec<Value>, RingBufferMetadata>,
) -> Result<()> {
if !cache.contains_key(partition_key) {
let existing =
catalog.find_partition_metadata(&mut Transaction::Command(txn), ringbuffer, partition_key)?;
let m = existing.unwrap_or_else(RingBufferMetadata::new);
cache.insert(partition_key.to_vec(), m);
}
Ok(())
}
fn evict_oldest_for_partition(
txn: &mut CommandTransaction,
ringbuffer: &RingBuffer,
partition: Option<Partition>,
metadata: &mut RingBufferMetadata,
) -> Result<()> {
if let Some(partition) = partition {
let range = PartitionedRowKey::partition_scan_range(ringbuffer.id, partition, None);
let oldest = txn.range_rev(range, RangeScope::All, 1)?.next().transpose()?;
if let Some(entry) = oldest
&& let TaggedKey::PartitionedRow(pk) = &entry.key
{
txn.remove_from_ringbuffer(ringbuffer, Some(partition), pk.row)?;
}
metadata.count -= 1;
return Ok(());
}
let mut evict_pos = metadata.head;
loop {
let key = RowKey::new(ringbuffer.id, RowNumber(evict_pos));
if txn.get(&key)?.is_some() {
txn.remove_from_ringbuffer(ringbuffer, None, RowNumber(evict_pos))?;
break;
}
evict_pos += 1;
if evict_pos >= metadata.tail {
break;
}
}
metadata.head = evict_pos + 1;
while metadata.head < metadata.tail {
let key = RowKey::new(ringbuffer.id, RowNumber(metadata.head));
if txn.get(&key)?.is_some() {
break;
}
metadata.head += 1;
}
metadata.count -= 1;
Ok(())
}
#[inline]
fn dict_encode_ringbuffer_row(
catalog: &Catalog,
txn: &mut CommandTransaction,
ringbuffer: &RingBuffer,
values: &mut [Value],
) -> Result<()> {
for (idx, col) in ringbuffer.columns.iter().enumerate() {
if let Some(dict_id) = col.dictionary_id {
let dictionary =
catalog.find_dictionary(&mut Transaction::Command(txn), dict_id)?.ok_or_else(|| {
internal_error!("Dictionary {:?} not found for column {}", dict_id, col.name)
})?;
let entry_id = if matches!(values[idx], Value::None { .. }) {
dictionary.id_type.none()
} else {
txn.insert_into_dictionary(&dictionary, &values[idx])?
};
values[idx] = entry_id.to_value();
}
}
Ok(())
}
#[inline]
fn dict_encode_series_row(
catalog: &Catalog,
txn: &mut CommandTransaction,
series: &Series,
values: &mut [Value],
) -> Result<()> {
for (idx, col) in series.columns.iter().enumerate() {
if let Some(dict_id) = col.dictionary_id {
let dictionary =
catalog.find_dictionary(&mut Transaction::Command(txn), dict_id)?.ok_or_else(|| {
internal_error!("Dictionary {:?} not found for column {}", dict_id, col.name)
})?;
let entry_id = if matches!(values[idx], Value::None { .. }) {
dictionary.id_type.none()
} else {
txn.insert_into_dictionary(&dictionary, &values[idx])?
};
values[idx] = entry_id.to_value();
}
}
Ok(())
}
fn execute_series_insert<V: ValidationMode>(
catalog: &Catalog,
txn: &mut CommandTransaction,
pending: &PendingSeriesInsert,
clock: &Clock,
) -> Result<SeriesInsertResult> {
let series = resolve_series(catalog, txn, pending)?;
let mut metadata = load_series_metadata(catalog, txn, pending, &series)?;
let shape = get_or_create_series_shape(catalog, &series, &mut Transaction::Command(txn))?;
let coerced_rows = coerce_series_rows::<V>(pending, &series, txn.identity)?;
let inserted = insert_series_rows::<V>(catalog, txn, &series, &shape, coerced_rows, &mut metadata, clock)?;
catalog.update_series_metadata_txn(&mut Transaction::Command(txn), series.id, metadata)?;
Ok(SeriesInsertResult {
namespace: pending.namespace.clone(),
series: pending.series.clone(),
inserted,
})
}
#[inline]
fn resolve_series(catalog: &Catalog, txn: &mut CommandTransaction, pending: &PendingSeriesInsert) -> Result<Series> {
let namespace = catalog
.find_namespace_by_name(&mut Transaction::Command(txn), &pending.namespace)?
.ok_or_else(|| CatalogError::NotFound {
kind: CatalogObjectKind::Namespace,
namespace: pending.namespace.to_string(),
name: String::new(),
fragment: Fragment::None,
})?;
catalog.find_series_by_name(&mut Transaction::Command(txn), namespace.id(), &pending.series)?.ok_or_else(|| {
CatalogError::NotFound {
kind: CatalogObjectKind::Series,
namespace: pending.namespace.to_string(),
name: pending.series.to_string(),
fragment: Fragment::None,
}
.into()
})
}
#[inline]
fn load_series_metadata(
catalog: &Catalog,
txn: &mut CommandTransaction,
pending: &PendingSeriesInsert,
series: &Series,
) -> Result<SeriesMetadata> {
catalog.find_series_metadata(&mut Transaction::Command(txn), series.id)?.ok_or_else(|| {
CatalogError::NotFound {
kind: CatalogObjectKind::Series,
namespace: pending.namespace.to_string(),
name: pending.series.to_string(),
fragment: Fragment::None,
}
.into()
})
}
#[inline]
fn coerce_series_rows<V: ValidationMode>(
pending: &PendingSeriesInsert,
series: &Series,
identity: IdentityId,
) -> Result<Vec<Vec<Value>>> {
if V::VALIDATED {
validate_and_coerce_rows_series(&pending.rows, series, identity)
} else {
reorder_rows_unvalidated_series(&pending.rows, series)
}
}
fn insert_series_rows<V: ValidationMode>(
catalog: &Catalog,
txn: &mut CommandTransaction,
series: &Series,
shape: &RowShape,
coerced_rows: Vec<Vec<Value>>,
metadata: &mut SeriesMetadata,
clock: &Clock,
) -> Result<u64> {
let key_col_name = series.key.column();
let key_col_idx =
series.columns.iter().position(|c| c.name == key_col_name).ok_or_else(|| {
internal_error!("series {} key column {} not found", series.name, key_col_name)
})?;
let partition_col_indices = series_partition_col_indices(series)?;
let storage = StorageId::series(series.id);
let mut verified: HashSet<Partition> = HashSet::new();
let mut inserted_count = 0u64;
for mut values in coerced_rows {
dict_encode_series_row(catalog, txn, series, &mut values)?;
if V::VALIDATED {
for (idx, col) in series.columns.iter().enumerate() {
col.constraint.validate(&values[idx])?;
}
}
let key_value = series.key_to_u64(values[key_col_idx].clone()).unwrap_or(0);
metadata.sequence_counter += 1;
let sequence = metadata.sequence_counter;
let key: TaggedKey = if partition_col_indices.is_empty() {
SeriesRowKey {
storage,
variant_tag: None,
key: key_value,
sequence,
}
.into()
} else {
let partition_values: Vec<Value> =
partition_col_indices.iter().map(|&idx| values[idx].clone()).collect();
let partition = Partition::of(&partition_values);
resolve_partition(
&mut Transaction::Command(txn),
ObjectId::Series(series.id),
partition,
&partition_values,
&mut verified,
)?;
PartitionedSeriesRowKey::new(storage, partition, None, key_value, sequence).into()
};
let row = encode_series_row(series, shape, key_value, &values, key_col_idx, clock)?;
let mut rows_buf = [row];
SeriesRowInterceptor::pre_insert(txn, series, &mut rows_buf)?;
let [row] = rows_buf;
let row = row.freeze_bytes();
txn.set(&key, row.clone())?;
let rows = [row.clone()];
SeriesRowInterceptor::post_insert(txn, series, &rows)?;
update_series_metadata_for_insert(metadata, key_value);
inserted_count += 1;
}
Ok(inserted_count)
}
#[inline]
fn series_partition_col_indices(series: &Series) -> Result<Vec<usize>> {
series.partition_by
.iter()
.map(|name| {
series.columns.iter().position(|c| &c.name == name).ok_or_else(|| {
internal_error!("series {} partition column {} not found", series.name, name)
})
})
.collect()
}
#[inline]
fn encode_series_row(
series: &Series,
shape: &RowShape,
key_value: u64,
values: &[Value],
key_col_idx: usize,
clock: &Clock,
) -> Result<EncodedSeriesRowBuilder> {
let key_value_encoded = series.key_from_u64(key_value);
let mut row = shape.allocate_series();
shape.set_value(&mut row, 0, &key_value_encoded);
let mut shape_idx = 1;
for (col_idx, value) in values.iter().enumerate() {
if col_idx == key_col_idx {
continue;
}
shape.set_value(&mut row, shape_idx, value);
shape_idx += 1;
}
let now = clock.now();
row.set_timestamps(now, now);
if let Some(time) = resolve_time(&series.name, &series.columns, &series.time, shape, &row, now)? {
row.set_time(time);
}
Ok(row)
}
#[inline]
fn update_series_metadata_for_insert(metadata: &mut SeriesMetadata, key_value: u64) {
if metadata.row_count == 0 {
metadata.oldest_key = key_value;
metadata.newest_key = key_value;
} else {
if key_value < metadata.oldest_key {
metadata.oldest_key = key_value;
}
if key_value > metadata.newest_key {
metadata.newest_key = key_value;
}
}
metadata.row_count += 1;
}
fn parse_qualified_name(qualified_name: &str) -> (String, String) {
if let Some((ns, name)) = qualified_name.rsplit_once("::") {
(ns.to_string(), name.to_string())
} else {
("default".to_string(), qualified_name.to_string())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parse_qualified_name_simple() {
assert_eq!(parse_qualified_name("table"), ("default".to_string(), "table".to_string()));
}
#[test]
fn parse_qualified_name_single_namespace() {
assert_eq!(parse_qualified_name("ns::table"), ("ns".to_string(), "table".to_string()));
}
#[test]
fn parse_qualified_name_nested_namespace() {
assert_eq!(parse_qualified_name("a::b::table"), ("a::b".to_string(), "table".to_string()));
}
#[test]
fn parse_qualified_name_deeply_nested_namespace() {
assert_eq!(parse_qualified_name("a::b::c::table"), ("a::b::c".to_string(), "table".to_string()));
}
#[test]
fn parse_qualified_name_empty_string() {
assert_eq!(parse_qualified_name(""), ("default".to_string(), "".to_string()));
}
}