use std::{
collections::{HashMap, HashSet},
iter,
};
use reifydb_codec::key::{encoded::EncodedKey, serializer::KeySerializer};
use reifydb_core::{interface::catalog::flow::FlowNodeId, key::flow_node_internal_state::FlowNodeInternalStateKey};
use reifydb_value::value::row_number::RowNumber;
use crate::{
error::Result,
operator::context::{InternalStateApi, OperatorContext},
};
pub struct RowNumberProvider {
_node: FlowNodeId,
}
impl RowNumberProvider {
pub fn new(node: FlowNodeId) -> Self {
Self {
_node: node,
}
}
pub fn get_or_create_row_numbers_batch<'a, O, I>(&self, ctx: &mut O, keys: I) -> Result<Vec<(RowNumber, bool)>>
where
O: OperatorContext,
I: IntoIterator<Item = &'a EncodedKey>,
{
let map_keys: Vec<EncodedKey> = keys.into_iter().map(|key| self.make_map_key(key)).collect();
let mut existing: HashMap<Vec<u8>, u64> = HashMap::with_capacity(map_keys.len());
ctx.internal_state().get_many_visit::<u64>(&map_keys, &mut |map_key, row_num| {
existing.insert(map_key.as_bytes().to_vec(), row_num);
Ok(())
})?;
let mut distinct_new: HashSet<Vec<u8>> = HashSet::new();
for map_key in &map_keys {
let bytes = map_key.as_bytes();
if !existing.contains_key(bytes) {
distinct_new.insert(bytes.to_vec());
}
}
let mut next = if distinct_new.is_empty() {
0
} else {
ctx.allocate_row_numbers(distinct_new.len() as u64)?.0
};
let mut newly_assigned: HashMap<Vec<u8>, u64> = HashMap::new();
let mut results = Vec::with_capacity(map_keys.len());
for map_key in &map_keys {
let bytes = map_key.as_bytes();
if let Some(&row_num) = existing.get(bytes).or_else(|| newly_assigned.get(bytes)) {
results.push((RowNumber(row_num), false));
continue;
}
let row_num = next;
next += 1;
ctx.internal_state().set::<u64>(map_key, &row_num)?;
newly_assigned.insert(bytes.to_vec(), row_num);
results.push((RowNumber(row_num), true));
}
Ok(results)
}
pub fn get_or_create_row_number<O: OperatorContext>(
&self,
ctx: &mut O,
key: &EncodedKey,
) -> Result<(RowNumber, bool)> {
Ok(self.get_or_create_row_numbers_batch(ctx, iter::once(key))?.into_iter().next().unwrap())
}
fn make_map_key(&self, key: &EncodedKey) -> EncodedKey {
let mut serializer = KeySerializer::new();
serializer.extend_u8(FlowNodeInternalStateKey::ROW_NUMBER_MAPPING_TAG);
serializer.extend_bytes(key.as_ref());
serializer.finish()
}
}
#[cfg(test)]
pub mod tests {
use reifydb_abi::operator::capabilities::OperatorCapability;
use reifydb_codec::key::encoded::EncodedKey;
use reifydb_core::interface::catalog::flow::FlowNodeId;
use crate::{
config::Config,
error::Result,
operator::{
FFIOperator, OperatorMetadata, change::BorrowedChange, column::operator::OperatorColumn,
context::ffi::FFIOperatorContext,
},
state::{RawStatefulOperator, row::RowNumberProvider},
testing::{harness::FFIOperatorHarnessBuilder, helpers::encode_key},
};
struct RowNumberTestOperator;
impl OperatorMetadata for RowNumberTestOperator {
const NAME: &'static str = "row_number_test";
const API: u32 = 1;
const VERSION: &'static str = "1.0.0";
const DESCRIPTION: &'static str = "Test operator for row number provider";
const INPUT_COLUMNS: &'static [OperatorColumn] = &[];
const OUTPUT_COLUMNS: &'static [OperatorColumn] = &[];
const CAPABILITIES: &'static [OperatorCapability] = OperatorCapability::STANDARD;
}
impl FFIOperator for RowNumberTestOperator {
fn new(_operator_id: FlowNodeId, _config: &Config) -> Result<Self> {
Ok(Self)
}
fn apply(&mut self, _ctx: &mut FFIOperatorContext, _input: BorrowedChange<'_>) -> Result<()> {
Ok(())
}
}
impl RawStatefulOperator for RowNumberTestOperator {}
#[test]
fn test_first_row_number_starts_at_one() {
let mut harness = FFIOperatorHarnessBuilder::<RowNumberTestOperator>::new()
.with_node_id(FlowNodeId(1))
.build()
.expect("Failed to build harness");
let key = encode_key("test_key");
let mut ctx = harness.create_operator_context();
let (row_num, is_new) = ctx.get_or_create_row_number(&key).unwrap();
assert_eq!(row_num.0, 1);
assert!(is_new);
}
#[test]
fn test_duplicate_key_returns_same_row_number() {
let mut harness = FFIOperatorHarnessBuilder::<RowNumberTestOperator>::new()
.with_node_id(FlowNodeId(1))
.build()
.expect("Failed to build harness");
let key = encode_key("test_key");
let mut ctx = harness.create_operator_context();
let (row_num1, is_new1) = ctx.get_or_create_row_number(&key).unwrap();
let mut ctx = harness.create_operator_context();
let (row_num2, is_new2) = ctx.get_or_create_row_number(&key).unwrap();
assert_eq!(row_num1.0, row_num2.0);
assert!(is_new1);
assert!(!is_new2);
}
#[test]
fn test_sequential_numbering() {
let mut harness = FFIOperatorHarnessBuilder::<RowNumberTestOperator>::new()
.with_node_id(FlowNodeId(1))
.build()
.expect("Failed to build harness");
let key1 = encode_key("key1");
let key2 = encode_key("key2");
let key3 = encode_key("key3");
let mut ctx = harness.create_operator_context();
let (row_num1, _) = ctx.get_or_create_row_number(&key1).unwrap();
let mut ctx = harness.create_operator_context();
let (row_num2, _) = ctx.get_or_create_row_number(&key2).unwrap();
let mut ctx = harness.create_operator_context();
let (row_num3, _) = ctx.get_or_create_row_number(&key3).unwrap();
assert_eq!(row_num1.0, 1);
assert_eq!(row_num2.0, 2);
assert_eq!(row_num3.0, 3);
}
#[test]
fn test_operator_isolation() {
let mut harness1 = FFIOperatorHarnessBuilder::<RowNumberTestOperator>::new()
.with_node_id(FlowNodeId(1))
.build()
.expect("Failed to build harness1");
let mut harness2 = FFIOperatorHarnessBuilder::<RowNumberTestOperator>::new()
.with_node_id(FlowNodeId(2))
.build()
.expect("Failed to build harness2");
let key = encode_key("same_key");
let mut ctx1 = harness1.create_operator_context();
let (row_num1, is_new1) = ctx1.get_or_create_row_number(&key).unwrap();
let mut ctx2 = harness2.create_operator_context();
let (row_num2, is_new2) = ctx2.get_or_create_row_number(&key).unwrap();
assert!(is_new1);
assert!(is_new2);
assert_eq!(row_num1.0, 1);
assert_eq!(row_num2.0, 1);
}
#[test]
fn test_persistence_across_calls() {
let mut harness = FFIOperatorHarnessBuilder::<RowNumberTestOperator>::new()
.with_node_id(FlowNodeId(1))
.build()
.expect("Failed to build harness");
let key1 = encode_key("key1");
let key2 = encode_key("key2");
let mut ctx = harness.create_operator_context();
ctx.get_or_create_row_number(&key1).unwrap();
let mut ctx = harness.create_operator_context();
ctx.get_or_create_row_number(&key2).unwrap();
let key3 = encode_key("key3");
let mut ctx = harness.create_operator_context();
let (row_num3, is_new3) = ctx.get_or_create_row_number(&key3).unwrap();
assert!(is_new3);
assert_eq!(row_num3.0, 3);
let mut ctx = harness.create_operator_context();
let (row_num1, is_new1) = ctx.get_or_create_row_number(&key1).unwrap();
assert!(!is_new1);
assert_eq!(row_num1.0, 1);
}
#[test]
fn test_large_scale_row_numbers() {
let mut harness = FFIOperatorHarnessBuilder::<RowNumberTestOperator>::new()
.with_node_id(FlowNodeId(1))
.build()
.expect("Failed to build harness");
for i in 0..1000 {
let key = encode_key(format!("key_{}", i));
let mut ctx = harness.create_operator_context();
let (row_num, is_new) = ctx.get_or_create_row_number(&key).unwrap();
assert!(is_new);
assert_eq!(row_num.0, i + 1);
}
let key_500 = encode_key("key_500");
let mut ctx = harness.create_operator_context();
let (row_num, is_new) = ctx.get_or_create_row_number(&key_500).unwrap();
assert!(!is_new);
assert_eq!(row_num.0, 501);
}
#[test]
fn test_empty_key() {
let mut harness = FFIOperatorHarnessBuilder::<RowNumberTestOperator>::new()
.with_node_id(FlowNodeId(1))
.build()
.expect("Failed to build harness");
let empty_key = encode_key("");
let mut ctx = harness.create_operator_context();
let (row_num, is_new) = ctx.get_or_create_row_number(&empty_key).unwrap();
assert!(is_new);
assert_eq!(row_num.0, 1);
let mut ctx = harness.create_operator_context();
let (row_num2, is_new2) = ctx.get_or_create_row_number(&empty_key).unwrap();
assert!(!is_new2);
assert_eq!(row_num2.0, 1);
}
#[test]
fn test_binary_key_data() {
let mut harness = FFIOperatorHarnessBuilder::<RowNumberTestOperator>::new()
.with_node_id(FlowNodeId(1))
.build()
.expect("Failed to build harness");
let binary_key = EncodedKey::new(vec![0x00, 0xFF, 0x00, 0xAB, 0xCD]);
let mut ctx = harness.create_operator_context();
let (row_num, is_new) = ctx.get_or_create_row_number(&binary_key).unwrap();
assert!(is_new);
assert_eq!(row_num.0, 1);
let mut ctx = harness.create_operator_context();
let (row_num2, is_new2) = ctx.get_or_create_row_number(&binary_key).unwrap();
assert!(!is_new2);
assert_eq!(row_num2.0, 1);
}
#[test]
fn test_interleaved_operations() {
let mut harness = FFIOperatorHarnessBuilder::<RowNumberTestOperator>::new()
.with_node_id(FlowNodeId(1))
.build()
.expect("Failed to build harness");
let key1 = encode_key("key1");
let key2 = encode_key("key2");
let mut ctx = harness.create_operator_context();
let (row_num1_first, _) = ctx.get_or_create_row_number(&key1).unwrap();
assert_eq!(row_num1_first.0, 1);
let mut ctx = harness.create_operator_context();
let (row_num2_first, _) = ctx.get_or_create_row_number(&key2).unwrap();
assert_eq!(row_num2_first.0, 2);
let mut ctx = harness.create_operator_context();
let (row_num1_second, is_new1) = ctx.get_or_create_row_number(&key1).unwrap();
assert!(!is_new1);
assert_eq!(row_num1_second.0, 1);
let mut ctx = harness.create_operator_context();
let (row_num2_second, is_new2) = ctx.get_or_create_row_number(&key2).unwrap();
assert!(!is_new2);
assert_eq!(row_num2_second.0, 2);
}
#[test]
fn test_map_key_uniqueness() {
let provider = RowNumberProvider::new(FlowNodeId(42));
let original_key1 = encode_key("test1");
let original_key2 = encode_key("test2");
let map_key1 = provider.make_map_key(&original_key1);
let map_key2 = provider.make_map_key(&original_key2);
assert!(!map_key1.is_empty());
assert!(!map_key2.is_empty());
assert_ne!(map_key1, map_key2);
let map_key1_again = provider.make_map_key(&original_key1);
assert_eq!(map_key1, map_key1_again);
}
#[test]
fn test_batch_mixed_existing_and_new_keys() {
let mut harness = FFIOperatorHarnessBuilder::<RowNumberTestOperator>::new()
.with_node_id(FlowNodeId(1))
.build()
.expect("Failed to build harness");
let provider = RowNumberProvider::new(FlowNodeId(1));
let key1 = encode_key("batch_key_1");
let key2 = encode_key("batch_key_2");
let key3 = encode_key("batch_key_3");
let mut ctx = harness.create_operator_context();
let (rn1, _) = provider.get_or_create_row_number(&mut ctx, &key1).unwrap();
assert_eq!(rn1.0, 1);
let mut ctx = harness.create_operator_context();
let (rn2, _) = provider.get_or_create_row_number(&mut ctx, &key2).unwrap();
assert_eq!(rn2.0, 2);
let mut ctx = harness.create_operator_context();
let (rn3, _) = provider.get_or_create_row_number(&mut ctx, &key3).unwrap();
assert_eq!(rn3.0, 3);
let key4 = encode_key("batch_key_4");
let key5 = encode_key("batch_key_5");
let batch_keys = vec![&key2, &key4, &key1, &key5, &key3];
let mut ctx = harness.create_operator_context();
let results = provider.get_or_create_row_numbers_batch(&mut ctx, batch_keys.into_iter()).unwrap();
assert_eq!(results.len(), 5);
assert_eq!(results[0].0.0, 2);
assert!(!results[0].1);
assert_eq!(results[1].0.0, 4);
assert!(results[1].1);
assert_eq!(results[2].0.0, 1);
assert!(!results[2].1);
assert_eq!(results[3].0.0, 5);
assert!(results[3].1);
assert_eq!(results[4].0.0, 3);
assert!(!results[4].1);
let key6 = encode_key("batch_key_6");
let mut ctx = harness.create_operator_context();
let (rn6, is_new6) = provider.get_or_create_row_number(&mut ctx, &key6).unwrap();
assert_eq!(rn6.0, 6);
assert!(is_new6);
let mut ctx = harness.create_operator_context();
let (check_rn4, is_new4) = provider.get_or_create_row_number(&mut ctx, &key4).unwrap();
assert_eq!(check_rn4.0, 4);
assert!(!is_new4);
let mut ctx = harness.create_operator_context();
let (check_rn5, is_new5) = provider.get_or_create_row_number(&mut ctx, &key5).unwrap();
assert_eq!(check_rn5.0, 5);
assert!(!is_new5);
}
}