use std::collections::HashMap;
use reifydb_codec::{
key::{encoded::EncodedKey, serializer::KeySerializer},
row::{
bytes::{EncodedBytes, RowBuilder, SHAPE_HEADER_SIZE, read_created_at},
shape::{RowFamily, RowShape},
table::EncodedTableRow,
},
};
use reifydb_core::{
interface::{
catalog::{
dictionary::Dictionary,
flow::OperatorId,
object::ObjectId,
storage::StorageId,
view::{View, ViewSortKey},
},
change::{Change, Diff},
flow::OperatorCapability,
resolved::ResolvedView,
},
key::{
row::{PartitionedRowKey, PartitionedSortedViewRowKey, RowKey, SortedViewRowKey},
sort_run::SortRun,
},
partition::partition_col_indices,
row::row_shape_from_columns,
value::column::{buffer::ColumnBuffer, columns::Columns},
};
use reifydb_transaction::interceptor::dictionary_row::DictionaryRowInterceptor;
use reifydb_value::{
Result,
error::Error,
value::{Value, datetime::DateTime, partition::Partition, row_number::RowNumber, value_type::ValueType},
};
use tracing::instrument;
use super::{
DurableSink, coerce_columns, emit_view_change, encode_row_at_index,
partition::{ensure_partition_unchanged, partition_of, resolve_partition_flow},
shape_field_columns,
};
use crate::{
error::FlowSinkError,
transaction::{FlowTransaction, deferred::DeferredTransaction},
};
const CREATED_AT_CACHE_CAPACITY: usize = 16_384;
pub struct SinkTableViewOperator {
operator: OperatorId,
view: ResolvedView,
storage: StorageId,
shape: RowShape,
sort: Vec<ViewSortKey>,
partition_indices: Vec<usize>,
verified_partitions: HashMap<Partition, Vec<Value>>,
created_at: HashMap<RowNumber, DateTime>,
}
impl SinkTableViewOperator {
pub fn new(operator: OperatorId, view: ResolvedView, partition_by: Vec<String>) -> Self {
let storage = view.def().storage_id();
let shape = row_shape_from_columns(RowFamily::Table, view.def().columns());
let sort = view.def().sort().to_vec();
let partition_indices = partition_col_indices(view.def().columns(), &partition_by);
Self {
operator,
view,
storage,
shape,
sort,
partition_indices,
verified_partitions: HashMap::new(),
created_at: HashMap::new(),
}
}
#[inline]
fn is_partitioned(&self) -> bool {
!self.partition_indices.is_empty()
}
#[inline]
fn row_key(&self, row: RowNumber) -> EncodedKey {
RowKey::encoded(self.storage, row)
}
#[inline]
fn sort_run(&self, cols: &Columns, row_idx: usize) -> SortRun {
let mut serializer = KeySerializer::new();
for key in &self.sort {
let value = cols.data_at(key.column.0 as usize).get_value(row_idx);
serializer.extend_value_with_direction(&value, key.direction.clone().into());
}
SortRun::from_encoded(serializer.to_encoded_key())
}
#[inline]
fn sorted_view_key(&self, cols: &Columns, row_idx: usize, row: RowNumber) -> EncodedKey {
if self.sort.is_empty() {
return self.row_key(row);
}
SortedViewRowKey::encoded(self.storage, self.sort_run(cols, row_idx), row)
}
#[inline]
fn partitioned_key(&self, cols: &Columns, row_idx: usize, partition: Partition, row: RowNumber) -> EncodedKey {
if self.sort.is_empty() {
return PartitionedRowKey::encoded(self.storage, partition, row);
}
PartitionedSortedViewRowKey::encoded(self.storage, partition, self.sort_run(cols, row_idx), row)
}
}
impl DurableSink for SinkTableViewOperator {
fn id(&self) -> OperatorId {
self.operator
}
fn capabilities(&self) -> &[OperatorCapability] {
OperatorCapability::STANDARD
}
fn apply(&mut self, txn: &mut DeferredTransaction, change: Change) -> Result<Change> {
for diff in change.diffs.iter() {
match diff {
Diff::Insert {
post,
..
} => self.apply_table_view_insert(txn, post)?,
Diff::Update {
pre,
post,
..
} => self.apply_table_view_update(txn, pre, post)?,
Diff::Remove {
pre,
..
} => self.apply_table_view_remove(txn, pre)?,
}
}
Ok(Change::from_flow(self.operator, change.version, Vec::new(), change.changed_at))
}
}
impl SinkTableViewOperator {
#[inline]
#[instrument(name = "flow::operator::sink::view::insert", level = "trace", skip_all, fields(rows = post.row_count()))]
fn apply_table_view_insert(&mut self, txn: &mut DeferredTransaction, post: &Columns) -> Result<()> {
let coerced = coerce_columns(post, self.view.def().columns())?;
let dict_encoded = dictionary_encode_view_columns(txn, self.view.def(), &coerced)?;
let source = dict_encoded.as_ref().unwrap_or(&coerced);
let row_count = source.row_count();
let field_columns = shape_field_columns(source, &self.shape);
let mut keys: Vec<EncodedKey> = Vec::with_capacity(row_count);
let mut encoded_bytes_list: Vec<EncodedBytes> = Vec::with_capacity(row_count);
for row_idx in 0..row_count {
let row_number = source.row_numbers()[row_idx];
let (_, encoded) =
encode_row_at_index(source, row_idx, &self.shape, row_number, &field_columns)?;
let key = if self.is_partitioned() {
let (partition, values) = partition_of(&self.partition_indices, &coerced, row_idx);
resolve_partition_flow(
txn,
ObjectId::from(self.storage),
partition,
&values,
&mut self.verified_partitions,
)?;
self.partitioned_key(source, row_idx, partition, row_number)
} else {
self.sorted_view_key(source, row_idx, row_number)
};
remember_created_at(&mut self.created_at, row_number, read_created_at(&encoded));
keys.push(key);
encoded_bytes_list.push(encoded);
}
txn.set_batch(&keys, &encoded_bytes_list)?;
emit_view_change(txn, self.view.def(), Diff::insert(coerced));
Ok(())
}
#[inline]
#[instrument(name = "flow::operator::sink::view::update", level = "trace", skip_all, fields(rows = post.row_count()))]
fn apply_table_view_update(
&mut self,
txn: &mut DeferredTransaction,
pre: &Columns,
post: &Columns,
) -> Result<()> {
let coerced_pre = coerce_columns(pre, self.view.def().columns())?;
let coerced_post = coerce_columns(post, self.view.def().columns())?;
let dict_pre = dictionary_encode_view_columns(txn, self.view.def(), &coerced_pre)?;
let dict_post = dictionary_encode_view_columns(txn, self.view.def(), &coerced_post)?;
let source_pre = dict_pre.as_ref().unwrap_or(&coerced_pre);
let source_post = dict_post.as_ref().unwrap_or(&coerced_post);
let row_count = source_post.row_count();
let field_columns = shape_field_columns(source_post, &self.shape);
let mut pre_keys: Vec<EncodedKey> = Vec::with_capacity(row_count);
let mut post_keys: Vec<EncodedKey> = Vec::with_capacity(row_count);
let mut post_encoded_bytes_vec: Vec<EncodedBytes> = Vec::with_capacity(row_count);
for row_idx in 0..row_count {
let pre_row_number = source_pre.row_numbers()[row_idx];
let post_row_number = source_post.row_numbers()[row_idx];
let (_, mut post_encoded) = encode_row_at_index(
source_post,
row_idx,
&self.shape,
post_row_number,
&field_columns,
)?;
let (pre_key, post_key) = if self.is_partitioned() {
let (pre_partition, _pre_values) =
partition_of(&self.partition_indices, &coerced_pre, row_idx);
let (post_partition, post_values) =
partition_of(&self.partition_indices, &coerced_post, row_idx);
ensure_partition_unchanged(
ObjectId::from(self.storage),
pre_partition,
post_partition,
)?;
resolve_partition_flow(
txn,
ObjectId::from(self.storage),
post_partition,
&post_values,
&mut self.verified_partitions,
)?;
(
self.partitioned_key(source_pre, row_idx, pre_partition, pre_row_number),
self.partitioned_key(source_post, row_idx, post_partition, post_row_number),
)
} else {
(
self.sorted_view_key(source_pre, row_idx, pre_row_number),
self.sorted_view_key(source_post, row_idx, post_row_number),
)
};
let mut prior_created =
self.created_at.get(&post_row_number).copied().filter(|c| !c.is_epoch());
if prior_created.is_none() && pre_row_number != post_row_number {
prior_created = self.created_at.get(&pre_row_number).copied().filter(|c| !c.is_epoch());
}
if prior_created.is_none() {
prior_created = match txn.get(&post_key)? {
Some(prior) if prior.len() >= SHAPE_HEADER_SIZE => {
let c = read_created_at(&prior);
if !c.is_epoch() {
Some(c)
} else {
None
}
}
_ => None,
};
if prior_created.is_none() && pre_key.as_slice() != post_key.as_slice() {
prior_created = match txn.get(&pre_key)? {
Some(prior) if prior.len() >= SHAPE_HEADER_SIZE => {
let c = read_created_at(&prior);
if !c.is_epoch() {
Some(c)
} else {
None
}
}
_ => None,
};
}
}
if let Some(c) = prior_created
&& post_encoded.len() >= SHAPE_HEADER_SIZE
{
let updated = self.shape.updated_at(&post_encoded);
let mut builder = EncodedTableRow::from(post_encoded).thaw();
builder.set_timestamps(c, updated);
post_encoded = builder.freeze_bytes();
}
if pre_row_number != post_row_number {
self.created_at.remove(&pre_row_number);
}
remember_created_at(&mut self.created_at, post_row_number, read_created_at(&post_encoded));
pre_keys.push(pre_key);
post_keys.push(post_key);
post_encoded_bytes_vec.push(post_encoded);
}
txn.remove_batch(&pre_keys)?;
txn.set_batch(&post_keys, &post_encoded_bytes_vec)?;
emit_view_change(txn, self.view.def(), Diff::update(coerced_pre, coerced_post));
Ok(())
}
#[inline]
#[instrument(name = "flow::operator::sink::view::remove", level = "trace", skip_all, fields(rows = pre.row_count()))]
fn apply_table_view_remove(&mut self, txn: &mut DeferredTransaction, pre: &Columns) -> Result<()> {
let coerced = coerce_columns(pre, self.view.def().columns())?;
let dict_encoded = dictionary_encode_view_columns(txn, self.view.def(), &coerced)?;
let source = dict_encoded.as_ref().unwrap_or(&coerced);
let row_count = source.row_count();
let mut keys: Vec<EncodedKey> = Vec::with_capacity(row_count);
for row_idx in 0..row_count {
let row_number = source.row_numbers()[row_idx];
self.created_at.remove(&row_number);
let key = if self.is_partitioned() {
let (partition, _values) = partition_of(&self.partition_indices, &coerced, row_idx);
self.partitioned_key(source, row_idx, partition, row_number)
} else {
self.sorted_view_key(source, row_idx, row_number)
};
keys.push(key);
}
txn.remove_batch(&keys)?;
emit_view_change(txn, self.view.def(), Diff::remove(coerced));
Ok(())
}
}
fn remember_created_at(cache: &mut HashMap<RowNumber, DateTime>, row_number: RowNumber, created_at: DateTime) {
if created_at.is_epoch() {
return;
}
if cache.len() >= CREATED_AT_CACHE_CAPACITY {
cache.clear();
}
cache.insert(row_number, created_at);
}
#[inline]
pub(crate) fn dictionary_encode_view_columns(
txn: &mut DeferredTransaction,
view: &View,
columns: &Columns,
) -> Result<Option<Columns>> {
let mut dict_columns: Vec<(usize, Dictionary)> = Vec::new();
{
let catalog = txn.catalog();
for (pos, col) in view.columns().iter().enumerate() {
if let Some(dict_id) = col.dictionary_id {
let dictionary = catalog.cache().find_dictionary(dict_id).ok_or_else(|| {
Error::from(FlowSinkError::DictionaryNotFound {
dictionary_id: format!("{:?}", dict_id),
column: col.name.to_string(),
})
})?;
dict_columns.push((pos, dictionary));
}
}
}
if dict_columns.is_empty() {
return Ok(None);
}
let mut encoded = columns.clone();
for (col_pos, dictionary) in &dict_columns {
let row_count = encoded[*col_pos].len();
let mut values: Vec<Value> = Vec::with_capacity(row_count);
for row_idx in 0..row_count {
let mut values_buf = [encoded[*col_pos].get_value(row_idx)];
DictionaryRowInterceptor::pre_insert(txn, dictionary, &mut values_buf)?;
let [value] = values_buf;
values.push(value);
}
let registry = txn.dictionary_allocators();
let outcomes = registry.intern_batch(dictionary, &values)?;
let mut new_data = ColumnBuffer::with_capacity(ValueType::DictionaryId, row_count);
for outcome in &outcomes {
new_data.push_value(outcome.id.to_value());
}
encoded.columns[*col_pos] = new_data;
}
Ok(Some(encoded))
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use postcard::from_bytes;
use reifydb_core::{
actors::pending::PendingWrite,
common::CommitVersion,
interface::{
catalog::{
column::{Column as CatalogColumn, ColumnIndex},
id::{ColumnId, NamespaceId, ViewId},
namespace::Namespace,
view::{TableView, ViewKind},
},
resolved::ResolvedNamespace,
store::SingleVersionGet,
},
key::{any::TaggedKey, catalog::DictionaryEntryIndexKey},
value::column::ColumnWithName,
};
use reifydb_test_harness::engine::TestEngine;
use reifydb_transaction::dictionary::{DictionaryAllocatorRegistry, store::SingleDictionaryStore};
use reifydb_value::{
fragment::Fragment,
value::{
constraint::TypeConstraint, datetime::DateTime, identity::IdentityId, row_number::RowNumber,
system_columns::SystemColumns, value_type::ValueType,
},
};
use super::*;
use crate::transaction::mock::FlowTxn;
fn test_view_def() -> View {
View::Table(TableView {
id: ViewId(1),
namespace: NamespaceId(1),
name: "v".to_string(),
kind: ViewKind::Deferred,
columns: vec![CatalogColumn {
id: ColumnId(1),
name: "v".to_string(),
constraint: TypeConstraint::unconstrained(ValueType::Float8),
properties: vec![],
index: ColumnIndex(0),
auto_increment: false,
dictionary_id: None,
}],
primary_key: None,
partition_by: vec![],
sort: vec![],
})
}
fn test_sink() -> SinkTableViewOperator {
let resolved = ResolvedView::new(
Fragment::internal("v"),
ResolvedNamespace::new(Fragment::internal("system"), Namespace::system()),
test_view_def(),
);
SinkTableViewOperator::new(OperatorId(1), resolved, vec![])
}
fn one_row(v: f64, ts_nanos: u64) -> Columns {
Columns::with_system(
vec![ColumnWithName::new(Fragment::internal("v"), ColumnBuffer::float8([v]))],
SystemColumns::new(
vec![RowNumber(1)],
Vec::new(),
vec![DateTime::from_nanos(ts_nanos)],
vec![DateTime::from_nanos(ts_nanos)],
vec![DateTime::from_nanos(ts_nanos)],
),
)
}
fn commit_flow_pending(engine: &TestEngine, txn: &mut DeferredTransaction) {
let pending = txn.take_pending();
let mut cmd = engine.begin_admin(IdentityId::system()).unwrap();
for (key, pw) in pending.iter_sorted() {
let key = TaggedKey::decode(key).unwrap();
match pw {
PendingWrite::Set(v) => cmd.set(&key, v.clone()).unwrap(),
PendingWrite::Remove {
..
} => cmd.remove(&key).unwrap(),
};
}
cmd.commit().unwrap();
}
fn stored_view_bytes(engine: &TestEngine, sink: &SinkTableViewOperator, rn: u64) -> EncodedTableRow {
let key = RowKey::new(sink.storage, RowNumber(rn));
let query = engine.inner().multi().begin_query().unwrap();
EncodedTableRow::from(query.get(&key).unwrap().expect("the view row must exist").bytes().clone())
}
#[test]
fn update_preserves_created_at_from_the_operator_cache_and_falls_back_after_rebuild() {
let engine = TestEngine::new();
let mut sink = test_sink();
let mut txn = engine.flow_txn().clock_millis(0).deferred();
sink.apply(
&mut txn,
Change::from_flow(
OperatorId(1),
CommitVersion(1),
vec![Diff::insert(one_row(1.0, 1_000))],
DateTime::from_nanos(0),
),
)
.unwrap();
commit_flow_pending(&engine, &mut txn);
assert_eq!(stored_view_bytes(&engine, &sink, 1).created_at(), DateTime::from_nanos(1_000));
let mut txn = engine.flow_txn().clock_millis(0).deferred();
sink.apply(
&mut txn,
Change::from_flow(
OperatorId(1),
CommitVersion(2),
vec![Diff::update(one_row(1.0, 1_000), one_row(2.0, 5_000))],
DateTime::from_nanos(0),
),
)
.unwrap();
commit_flow_pending(&engine, &mut txn);
let stored = stored_view_bytes(&engine, &sink, 1);
assert_eq!(
stored.created_at(),
DateTime::from_nanos(1_000),
"created_at must survive the cached update"
);
assert_eq!(stored.updated_at(), DateTime::from_nanos(5_000), "updated_at must advance on every update");
let mut rebuilt = test_sink();
let mut txn = engine.flow_txn().clock_millis(0).deferred();
rebuilt.apply(
&mut txn,
Change::from_flow(
OperatorId(1),
CommitVersion(3),
vec![Diff::update(one_row(2.0, 5_000), one_row(3.0, 9_000))],
DateTime::from_nanos(0),
),
)
.unwrap();
commit_flow_pending(&engine, &mut txn);
let stored = stored_view_bytes(&engine, &rebuilt, 1);
assert_eq!(
stored.created_at(),
DateTime::from_nanos(1_000),
"created_at must survive the fallback path too"
);
assert_eq!(stored.updated_at(), DateTime::from_nanos(9_000));
}
#[test]
fn a_cold_registry_seeds_past_every_durable_id_and_never_clobbers() {
let t = TestEngine::new();
t.admin("CREATE NAMESPACE test");
t.admin("CREATE DICTIONARY test::syms FOR utf8 AS uint2");
let engine = t.inner();
let catalog = engine.catalog();
let namespace = catalog.cache().find_namespace_by_name("test").expect("namespace test");
let dictionary =
catalog.cache().find_dictionary_by_name(namespace.id(), "syms").expect("dictionary syms");
let intern = |value: &str| -> u128 {
let registry = DictionaryAllocatorRegistry::new(Arc::new(SingleDictionaryStore::new(
engine.single().clone(),
)));
registry.intern(&dictionary, &Value::Utf8(value.to_string())).unwrap().id.to_u128()
};
let sol_id = intern("sol");
let usdc_id = intern("usdc");
assert_ne!(sol_id, usdc_id, "distinct strings must intern to distinct ids");
let wsol_id = intern("wsol");
assert_ne!(wsol_id, sol_id, "wsol must not reuse sol's id (would overwrite sol's entry)");
assert_ne!(wsol_id, usdc_id, "wsol must not reuse usdc's id (would overwrite usdc's entry)");
let decode = |id: u128| -> String {
let key = DictionaryEntryIndexKey::encoded(dictionary.id, id);
let store = engine.single().read_store();
let row = SingleVersionGet::get(&store, &key).unwrap().expect("index entry present");
match from_bytes::<Value>(&row.bytes).unwrap() {
Value::Utf8(s) => s,
other => panic!("expected Utf8, got {:?}", other),
}
};
assert_eq!(decode(sol_id), "sol", "sol's dictionary entry was overwritten");
assert_eq!(decode(usdc_id), "usdc", "usdc's dictionary entry was overwritten");
assert_eq!(decode(wsol_id), "wsol");
}
}