use crate::database::mem_buffer::IndexMemBuffer;
use crate::index::btree_generic::{BTreeKey, GenericBTree, GenericBTreeConfig};
use crate::index::cached_index::CachedIndex;
use crate::types::{RowId, Value};
use crate::{Result, StorageError};
use parking_lot::{Mutex, RwLock};
use std::collections::HashSet;
use std::path::{Path, PathBuf};
use std::sync::Arc;
#[cfg(feature = "rayon")]
use rayon::prelude::*;
#[derive(Debug, Clone)]
pub struct ColumnValueIndexConfig {
pub max_page_size: usize,
pub cache_size: usize,
pub mem_buffer_size: usize,
pub drain_threshold: usize,
}
impl Default for ColumnValueIndexConfig {
fn default() -> Self {
Self {
max_page_size: 4096,
cache_size: 1024,
mem_buffer_size: 1024 * 1024, drain_threshold: 2,
}
}
}
const VALUE_DATA_SIZE: usize = 64;
const ROW_ID_SIZE: usize = 8;
const VALUE_LEN_SIZE: usize = 2;
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
struct IndexKey {
value_bytes: [u8; VALUE_DATA_SIZE],
row_id: RowId,
}
impl std::hash::Hash for IndexKey {
fn hash<H: std::hash::Hasher>(&self, state: &mut H) {
self.value_bytes.hash(state);
self.row_id.hash(state);
}
}
fn tombstone_key(key: &IndexKey) -> IndexKey {
IndexKey {
value_bytes: key.value_bytes,
row_id: key.row_id,
}
}
impl BTreeKey for IndexKey {
fn serialize(&self) -> Vec<u8> {
let key_size = Self::key_size();
let mut result = vec![0u8; key_size];
result[..VALUE_DATA_SIZE].copy_from_slice(&self.value_bytes[..]);
result[VALUE_DATA_SIZE..VALUE_DATA_SIZE + ROW_ID_SIZE]
.copy_from_slice(&self.row_id.to_be_bytes());
let vlen = VALUE_DATA_SIZE as u16;
result[VALUE_DATA_SIZE + ROW_ID_SIZE..VALUE_DATA_SIZE + ROW_ID_SIZE + VALUE_LEN_SIZE]
.copy_from_slice(&vlen.to_be_bytes());
result
}
fn deserialize(bytes: &[u8]) -> Result<Self> {
let key_size = Self::key_size();
if bytes.len() < key_size {
return Err(StorageError::Serialization(
"Invalid key: too short".to_string(),
));
}
let mut value_bytes = [0u8; VALUE_DATA_SIZE];
value_bytes.copy_from_slice(&bytes[..VALUE_DATA_SIZE]);
let row_id = u64::from_be_bytes(
bytes[VALUE_DATA_SIZE..VALUE_DATA_SIZE + ROW_ID_SIZE]
.try_into()
.map_err(|_| StorageError::Serialization("Invalid row_id".to_string()))?,
);
Ok(IndexKey {
value_bytes,
row_id,
})
}
fn key_size() -> usize {
VALUE_DATA_SIZE + ROW_ID_SIZE + VALUE_LEN_SIZE }
}
pub struct ColumnValueIndex {
_table_name: String,
column_name: String,
_storage_path: PathBuf,
btree: Arc<RwLock<GenericBTree<IndexKey>>>,
lru_cache: Arc<CachedIndex>,
mem_buffer: IndexMemBuffer<IndexKey, ()>,
tombstones: Mutex<HashSet<IndexKey>>,
pending_deletes: Mutex<Vec<IndexKey>>,
drain_lock: Mutex<()>,
drain_threshold: usize,
needs_rebuild: std::sync::atomic::AtomicBool,
}
impl ColumnValueIndex {
pub fn create<P: AsRef<Path>>(
path: P,
table_name: String,
column_name: String,
config: ColumnValueIndexConfig,
) -> Result<Self> {
let storage_path = path.as_ref().to_path_buf();
let btree_config = GenericBTreeConfig {
cache_size: config.cache_size,
unique_keys: false,
allow_updates: true,
immediate_sync: false,
};
let btree = GenericBTree::with_config(storage_path.clone(), btree_config)?;
Ok(Self {
_table_name: table_name,
column_name,
_storage_path: storage_path,
btree: Arc::new(RwLock::new(btree)),
lru_cache: Arc::new(CachedIndex::new(500)),
mem_buffer: IndexMemBuffer::new(config.mem_buffer_size),
tombstones: Mutex::new(HashSet::new()),
pending_deletes: Mutex::new(Vec::new()),
drain_lock: Mutex::new(()),
drain_threshold: config.drain_threshold,
needs_rebuild: std::sync::atomic::AtomicBool::new(true),
})
}
pub fn open<P: AsRef<Path>>(
path: P,
table_name: String,
column_name: String,
config: ColumnValueIndexConfig,
) -> Result<Self> {
let index = Self::create(path, table_name, column_name, config)?;
index
.needs_rebuild
.store(false, std::sync::atomic::Ordering::Relaxed);
Ok(index)
}
pub fn insert(&self, value: &Value, row_id: RowId) -> Result<()> {
let value_bytes = self.value_to_bytes(value)?;
let key = IndexKey {
value_bytes,
row_id,
};
let full = self
.mem_buffer
.insert(key.clone(), ())
.map_err(StorageError::InvalidData)?;
self.tombstones.lock().remove(&tombstone_key(&key));
if full {
if let Some(_guard) = self.drain_lock.try_lock() {
self.drain_immutable_to_btree()?;
}
}
self.lru_cache.try_invalidate(value);
Ok(())
}
pub fn bulk_insert_entry(&self, entries: &[(Value, RowId)]) -> Result<()> {
if entries.is_empty() {
return Ok(());
}
let keys: Vec<IndexKey> = entries
.iter()
.map(|(value, row_id)| {
let value_bytes = self.value_to_bytes(value).unwrap_or([0u8; 64]);
IndexKey {
value_bytes,
row_id: *row_id,
}
})
.collect();
self.bulk_load_or_insert(keys)
}
pub fn bulk_insert_raw(&self, entries: Vec<([u8; 64], RowId)>) -> Result<()> {
if entries.is_empty() {
return Ok(());
}
let keys: Vec<IndexKey> = entries
.into_iter()
.map(|(value_bytes, row_id)| IndexKey {
value_bytes,
row_id,
})
.collect();
self.bulk_load_or_insert(keys)
}
fn bulk_load_or_insert(&self, mut keys: Vec<IndexKey>) -> Result<()> {
#[cfg(feature = "rayon")]
{
keys.par_sort_unstable();
}
#[cfg(not(feature = "rayon"))]
{
keys.sort_unstable();
}
keys.dedup();
let mut btree = self.btree.write();
btree.bulk_load(keys)?;
Ok(())
}
pub fn update(&self, old_value: &Value, new_value: &Value, row_id: RowId) -> Result<()> {
let old_value_bytes = self.value_to_bytes(old_value)?;
let new_value_bytes = self.value_to_bytes(new_value)?;
let old_key = IndexKey {
value_bytes: old_value_bytes,
row_id,
};
let new_key = IndexKey {
value_bytes: new_value_bytes,
row_id,
};
self.mem_buffer.delete(&old_key);
let pending_len = if old_key != new_key {
let mut pending = self.pending_deletes.lock();
pending.push(old_key.clone());
pending.len()
} else {
0
};
{
let mut tombstones = self.tombstones.lock();
tombstones.remove(&tombstone_key(&new_key)); if old_key != new_key {
tombstones.insert(tombstone_key(&old_key)); }
}
let full = self
.mem_buffer
.insert(new_key.clone(), ())
.map_err(StorageError::InvalidData)?;
if full || pending_len > 10_000 {
if let Some(_guard) = self.drain_lock.try_lock() {
self.drain_immutable_to_btree()?;
}
}
self.lru_cache.try_invalidate(old_value);
self.lru_cache.try_invalidate(new_value);
Ok(())
}
pub fn batch_insert(&self, items: Vec<(Value, RowId)>) -> Result<()> {
if items.is_empty() {
return Ok(());
}
let mut keys: Vec<(IndexKey, Value)> = items
.into_iter()
.map(|(value, row_id)| {
let value_bytes = self.value_to_bytes(&value)?;
let key = IndexKey {
value_bytes,
row_id,
};
Ok((key, value))
})
.collect::<Result<Vec<_>>>()?;
keys.sort_by(|a, b| a.0.value_bytes.cmp(&b.0.value_bytes));
{
let mut tombstones = self.tombstones.lock();
for (key, _) in &keys {
tombstones.remove(&tombstone_key(key));
}
}
let buffer_entries: Vec<(IndexKey, ())> =
keys.iter().map(|(k, _)| (k.clone(), ())).collect();
let full = self
.mem_buffer
.batch_insert(buffer_entries)
.map_err(StorageError::InvalidData)?;
if full {
if let Some(_guard) = self.drain_lock.try_lock() {
self.drain_immutable_to_btree()?;
}
}
for (_, value) in &keys {
self.lru_cache.invalidate(value);
}
Ok(())
}
pub fn get_arc(&self, value: &Value) -> Result<Arc<Vec<RowId>>> {
if let Some(cached_ids) = self.lru_cache.get(value) {
return Ok(cached_ids);
}
self.lru_cache.record_miss();
let value_bytes = self.value_to_bytes(value)?;
let start_key = IndexKey {
value_bytes,
row_id: 0,
};
let end_key = IndexKey {
value_bytes,
row_id: RowId::MAX,
};
let tombstones = self.tombstones.lock();
let mut results: Vec<IndexKey> = Vec::new();
let mut seen = HashSet::new();
let buffer_results = self.mem_buffer.range(&start_key, &end_key);
for (key, _) in buffer_results {
if key.value_bytes == value_bytes
&& !tombstones.contains(&tombstone_key(&key))
&& seen.insert(key.row_id)
{
results.push(key);
}
}
{
let btree = self.btree.read();
let btree_results = btree.range(&start_key, &end_key)?;
for (key, _) in btree_results {
if key.value_bytes == value_bytes
&& !tombstones.contains(&tombstone_key(&key))
&& seen.insert(key.row_id)
{
results.push(key);
}
}
}
let row_ids: Vec<RowId> = results.into_iter().map(|key| key.row_id).collect();
let arc = Arc::new(row_ids);
if !arc.is_empty() {
self.lru_cache.put(value.clone(), (*arc).clone());
}
drop(tombstones);
Ok(arc)
}
pub fn get(&self, value: &Value) -> Result<Vec<RowId>> {
if let Some(cached_ids) = self.lru_cache.get(value) {
return Ok((*cached_ids).clone());
}
self.lru_cache.record_miss();
let value_bytes = self.value_to_bytes(value)?;
let start_key = IndexKey {
value_bytes,
row_id: 0,
};
let end_key = IndexKey {
value_bytes,
row_id: RowId::MAX,
};
let tombstones = self.tombstones.lock();
let mut results: Vec<IndexKey> = Vec::new();
let mut seen = HashSet::new();
let buffer_results = self.mem_buffer.range(&start_key, &end_key);
for (key, _) in buffer_results {
if key.value_bytes == value_bytes
&& !tombstones.contains(&tombstone_key(&key))
&& seen.insert(key.row_id)
{
results.push(key);
}
}
{
let btree = self.btree.read();
let btree_results = btree.range(&start_key, &end_key)?;
for (key, _) in btree_results {
if key.value_bytes == value_bytes
&& !tombstones.contains(&tombstone_key(&key))
&& seen.insert(key.row_id)
{
results.push(key);
}
}
}
let filtered: Vec<RowId> = results.into_iter().map(|key| key.row_id).collect();
if !filtered.is_empty() {
self.lru_cache.put(value.clone(), filtered.clone());
}
drop(tombstones);
Ok(filtered)
}
pub fn range(&self, start: &Value, end: &Value) -> Result<Vec<RowId>> {
let start_bytes = self.value_to_bytes(start)?;
let end_bytes = self.value_to_bytes(end)?;
let start_key = IndexKey {
value_bytes: start_bytes,
row_id: 0,
};
let end_key = IndexKey {
value_bytes: end_bytes,
row_id: RowId::MAX,
};
let tombstones = self.tombstones.lock();
let mut results: Vec<IndexKey> = Vec::with_capacity(64);
let mut seen: HashSet<u64> = HashSet::with_capacity(64);
let buffer_results = self.mem_buffer.range(&start_key, &end_key);
for (key, _) in buffer_results {
if !tombstones.contains(&tombstone_key(&key)) && seen.insert(key.row_id) {
results.push(key);
}
}
{
let btree = self.btree.read();
let btree_results = btree.range(&start_key, &end_key)?;
for (key, _) in btree_results {
if !tombstones.contains(&tombstone_key(&key)) && seen.insert(key.row_id) {
results.push(key);
}
}
}
drop(tombstones);
let row_ids: Vec<RowId> = results.into_iter().map(|key| key.row_id).collect();
Ok(row_ids)
}
pub fn scan_row_ids_with_limit(&self, limit: Option<usize>) -> Result<Vec<RowId>> {
let min_key = IndexKey {
value_bytes: [0u8; VALUE_DATA_SIZE],
row_id: 0,
};
let max_key = IndexKey {
value_bytes: [0xFFu8; VALUE_DATA_SIZE],
row_id: RowId::MAX,
};
let tombstones = self.tombstones.lock();
let mut results: Vec<IndexKey> = Vec::new();
let mut seen = HashSet::new();
let buffer_results = self.mem_buffer.range(&min_key, &max_key);
for (key, _) in buffer_results {
if !tombstones.contains(&tombstone_key(&key)) && seen.insert(key.row_id) {
results.push(key);
}
}
{
let btree = self.btree.read();
let all_entries = if let Some(limit_count) = limit {
btree.range_with_limit(&min_key, &max_key, limit_count)?
} else {
btree.range(&min_key, &max_key)?
};
for (key, _) in all_entries {
if !tombstones.contains(&tombstone_key(&key)) && seen.insert(key.row_id) {
results.push(key);
}
}
}
drop(tombstones);
let row_ids: Vec<RowId> = results.into_iter().map(|key| key.row_id).collect();
Ok(row_ids)
}
pub fn query_less_than(&self, upper_bound: &Value) -> Result<Vec<RowId>> {
let upper_bytes = self.value_to_bytes(upper_bound)?;
let start_key = IndexKey {
value_bytes: [0u8; VALUE_DATA_SIZE],
row_id: 0,
};
let end_key = IndexKey {
value_bytes: upper_bytes,
row_id: RowId::MAX,
};
let tombstones = self.tombstones.lock();
let mut results: Vec<IndexKey> = Vec::new();
let mut seen = HashSet::new();
let buffer_results = self.mem_buffer.range(&start_key, &end_key);
for (key, _) in buffer_results {
if key.value_bytes != upper_bytes
&& !tombstones.contains(&tombstone_key(&key))
&& seen.insert(key.row_id)
{
results.push(key);
}
}
{
let btree = self.btree.read();
let btree_results = btree.range(&start_key, &end_key)?;
for (key, _) in btree_results {
if key.value_bytes != upper_bytes
&& !tombstones.contains(&tombstone_key(&key))
&& seen.insert(key.row_id)
{
results.push(key);
}
}
}
drop(tombstones);
let row_ids: Vec<RowId> = results.into_iter().map(|key| key.row_id).collect();
Ok(row_ids)
}
pub fn query_greater_than(&self, lower_bound: &Value) -> Result<Vec<RowId>> {
let lower_bytes = self.value_to_bytes(lower_bound)?;
let start_key = IndexKey {
value_bytes: lower_bytes,
row_id: 0,
};
let end_key = IndexKey {
value_bytes: [0xFFu8; VALUE_DATA_SIZE],
row_id: RowId::MAX,
};
let tombstones = self.tombstones.lock();
let mut results: Vec<IndexKey> = Vec::new();
let mut seen = HashSet::new();
let buffer_results = self.mem_buffer.range(&start_key, &end_key);
for (key, _) in buffer_results {
if key.value_bytes != lower_bytes
&& !tombstones.contains(&tombstone_key(&key))
&& seen.insert(key.row_id)
{
results.push(key);
}
}
{
let btree = self.btree.read();
let btree_results = btree.range(&start_key, &end_key)?;
for (key, _) in btree_results {
if key.value_bytes != lower_bytes
&& !tombstones.contains(&tombstone_key(&key))
&& seen.insert(key.row_id)
{
results.push(key);
}
}
}
drop(tombstones);
let row_ids: Vec<RowId> = results.into_iter().map(|key| key.row_id).collect();
Ok(row_ids)
}
pub fn query_less_than_or_equal(&self, upper_bound: &Value) -> Result<Vec<RowId>> {
let upper_bytes = self.value_to_bytes(upper_bound)?;
let start_key = IndexKey {
value_bytes: [0u8; VALUE_DATA_SIZE],
row_id: 0,
};
let end_key = IndexKey {
value_bytes: upper_bytes,
row_id: RowId::MAX,
};
let tombstones = self.tombstones.lock();
let mut results: Vec<IndexKey> = Vec::new();
let mut seen = HashSet::new();
let buffer_results = self.mem_buffer.range(&start_key, &end_key);
for (key, _) in buffer_results {
if !tombstones.contains(&tombstone_key(&key)) && seen.insert(key.row_id) {
results.push(key);
}
}
{
let btree = self.btree.read();
let btree_results = btree.range(&start_key, &end_key)?;
for (key, _) in btree_results {
if !tombstones.contains(&tombstone_key(&key)) && seen.insert(key.row_id) {
results.push(key);
}
}
}
drop(tombstones);
let row_ids: Vec<RowId> = results.into_iter().map(|key| key.row_id).collect();
Ok(row_ids)
}
pub fn query_greater_than_or_equal(&self, lower_bound: &Value) -> Result<Vec<RowId>> {
let lower_bytes = self.value_to_bytes(lower_bound)?;
let start_key = IndexKey {
value_bytes: lower_bytes,
row_id: 0,
};
let end_key = IndexKey {
value_bytes: [0xFFu8; VALUE_DATA_SIZE],
row_id: RowId::MAX,
};
let tombstones = self.tombstones.lock();
let mut results: Vec<IndexKey> = Vec::new();
let mut seen = HashSet::new();
let buffer_results = self.mem_buffer.range(&start_key, &end_key);
for (key, _) in buffer_results {
if !tombstones.contains(&tombstone_key(&key)) && seen.insert(key.row_id) {
results.push(key);
}
}
{
let btree = self.btree.read();
let btree_results = btree.range(&start_key, &end_key)?;
for (key, _) in btree_results {
if !tombstones.contains(&tombstone_key(&key)) && seen.insert(key.row_id) {
results.push(key);
}
}
}
drop(tombstones);
let row_ids: Vec<RowId> = results.into_iter().map(|key| key.row_id).collect();
Ok(row_ids)
}
pub fn query_between(
&self,
lower_bound: &Value,
lower_inclusive: bool,
upper_bound: &Value,
upper_inclusive: bool,
) -> Result<Vec<RowId>> {
let lower_bytes = self.value_to_bytes(lower_bound)?;
let upper_bytes = self.value_to_bytes(upper_bound)?;
let start_key = IndexKey {
value_bytes: lower_bytes,
row_id: if lower_inclusive { 0 } else { RowId::MAX },
};
let end_key = IndexKey {
value_bytes: upper_bytes,
row_id: RowId::MAX,
};
let tombstones = self.tombstones.lock();
let mut results: Vec<IndexKey> = Vec::new();
let mut seen = HashSet::new();
let mut accept = |key: &IndexKey| -> bool {
if !lower_inclusive && key.value_bytes == start_key.value_bytes {
return false;
}
if !upper_inclusive && key.value_bytes == end_key.value_bytes {
return false;
}
!tombstones.contains(&tombstone_key(key)) && seen.insert(key.row_id)
};
let buffer_results = self.mem_buffer.range(&start_key, &end_key);
for (key, _) in buffer_results {
if accept(&key) {
results.push(key);
}
}
{
let btree = self.btree.read();
let btree_results = btree.range(&start_key, &end_key)?;
for (key, _) in btree_results {
if accept(&key) {
results.push(key);
}
}
}
drop(tombstones);
let row_ids: Vec<RowId> = results.into_iter().map(|key| key.row_id).collect();
Ok(row_ids)
}
pub fn delete(&self, value: &Value, row_id: RowId) -> Result<()> {
let value_bytes = self.value_to_bytes(value)?;
let key = IndexKey {
value_bytes,
row_id,
};
self.mem_buffer.delete(&key);
self.tombstones.lock().insert(tombstone_key(&key));
let mut btree = self.btree.write();
btree.delete(&key)?;
drop(btree);
self.lru_cache.invalidate(value);
Ok(())
}
pub fn delete_range(&self, start: &Value, end: &Value) -> Result<usize> {
let start_bytes = self.value_to_bytes(start)?;
let end_bytes = self.value_to_bytes(end)?;
let start_key = IndexKey {
value_bytes: start_bytes,
row_id: 0,
};
let end_key = IndexKey {
value_bytes: end_bytes,
row_id: RowId::MAX,
};
let mut deleted_count = 0;
let buffer_keys: Vec<IndexKey> = self
.mem_buffer
.range(&start_key, &end_key)
.into_iter()
.map(|(k, _)| k)
.collect();
let mut tombstones = self.tombstones.lock();
let mut btree = self.btree.write();
let btree_keys: Vec<IndexKey> = btree
.range(&start_key, &end_key)?
.into_iter()
.map(|(key, _)| key)
.collect();
for key in &btree_keys {
btree.delete(key)?;
tombstones.insert(tombstone_key(key));
deleted_count += 1;
}
drop(btree);
let mem_tombstone_keys: Vec<IndexKey> = buffer_keys.iter().map(tombstone_key).collect();
for tk in &mem_tombstone_keys {
tombstones.insert(tk.clone());
deleted_count += 1;
}
drop(tombstones);
for key in &buffer_keys {
self.mem_buffer.delete(key);
}
self.lru_cache.invalidate_range(start, end);
Ok(deleted_count)
}
pub fn flush(&self) -> Result<()> {
self.flush_buffer()?;
let mut btree = self.btree.write();
btree.flush()?;
Ok(())
}
fn drain_immutable_to_btree(&self) -> Result<()> {
self.drain_immutable_to_btree_impl(false)
}
fn drain_immutable_to_btree_impl(&self, force: bool) -> Result<()> {
if !force && self.mem_buffer.immutable_count() < self.drain_threshold {
return Ok(());
}
while self.mem_buffer.should_flush() {
if let Some(entries) = self.mem_buffer.flush().map_err(StorageError::InvalidData)? {
if !entries.is_empty() {
let tombstones = self.tombstones.lock();
let mut btree = self.btree.write();
for (key, _) in entries {
if !tombstones.contains(&tombstone_key(&key)) {
btree.insert(key, vec![])?;
}
}
}
} else {
break;
}
}
let deletes: Vec<IndexKey> = {
let mut pending = self.pending_deletes.lock();
std::mem::take(&mut *pending)
};
if !deletes.is_empty() {
let mut btree = self.btree.write();
for key in &deletes {
let _ = btree.delete(key);
}
}
Ok(())
}
pub fn flush_buffer(&self) -> Result<()> {
let entries = self.mem_buffer.drain();
let has_entries = !entries.is_empty();
let deletes: Vec<IndexKey> = {
let mut pending = self.pending_deletes.lock();
std::mem::take(&mut *pending)
};
let has_deletes = !deletes.is_empty();
if has_entries || has_deletes {
let tombstone_keys_to_clear: Vec<IndexKey> = {
let tombstones = self.tombstones.lock();
let mut btree = self.btree.write();
let mut keys_to_clear = Vec::new();
for (key, _) in &entries {
let tk = tombstone_key(key);
if tombstones.contains(&tk) {
keys_to_clear.push(tk);
} else {
btree.insert(key.clone(), vec![])?;
}
}
for key in &deletes {
let _ = btree.delete(key);
keys_to_clear.push(tombstone_key(key));
}
drop(btree);
drop(tombstones);
keys_to_clear
};
let mut tombstones = self.tombstones.lock();
for tk in &tombstone_keys_to_clear {
tombstones.remove(tk);
}
}
Ok(())
}
pub fn stats(&self) -> IndexStats {
let lru_stats = self.lru_cache.stats();
IndexStats {
cached_values: lru_stats.size,
total_row_ids: 0,
}
}
pub fn needs_rebuild(&self) -> bool {
self.needs_rebuild
.load(std::sync::atomic::Ordering::Acquire)
}
pub fn mark_rebuilt(&self) {
self.needs_rebuild
.store(false, std::sync::atomic::Ordering::Release);
}
pub fn entry_count(&self) -> usize {
let btree = self.btree.read();
btree.approximate_entry_count()
}
pub fn all_keys(&self, col_type: &crate::types::ColumnType) -> Result<Vec<Value>> {
let mut seen = std::collections::HashSet::new();
let mut keys = Vec::new();
for (idx_key, _) in self.mem_buffer.scan_all() {
if seen.insert(idx_key.value_bytes) {
keys.push(Self::bytes_to_value(&idx_key.value_bytes, col_type));
}
}
{
let min_key = IndexKey {
value_bytes: [0u8; VALUE_DATA_SIZE],
row_id: 0,
};
let max_key = IndexKey {
value_bytes: [0xFFu8; VALUE_DATA_SIZE],
row_id: u64::MAX,
};
let btree = self.btree.read();
if let Ok(entries) = btree.range(&min_key, &max_key) {
for (idx_key, _) in entries {
if seen.insert(idx_key.value_bytes) {
keys.push(Self::bytes_to_value(&idx_key.value_bytes, col_type));
}
}
}
}
Ok(keys)
}
fn bytes_to_value(bytes: &[u8; VALUE_DATA_SIZE], col_type: &crate::types::ColumnType) -> Value {
match col_type {
crate::types::ColumnType::Integer => {
let i = i64::from_be_bytes(bytes[..8].try_into().unwrap_or([0; 8]));
Value::Integer(i)
}
crate::types::ColumnType::Float => {
let sortable = u64::from_be_bytes(bytes[..8].try_into().unwrap_or([0; 8]));
let bits = if sortable & (1u64 << 63) != 0 {
!sortable } else {
sortable ^ (1u64 << 63) };
Value::Float(f64::from_bits(bits))
}
crate::types::ColumnType::Timestamp => {
let ts = i64::from_be_bytes(bytes[..8].try_into().unwrap_or([0; 8]));
Value::Timestamp(crate::types::Timestamp::from_micros(ts))
}
crate::types::ColumnType::Boolean => Value::Bool(bytes[0] != 0),
crate::types::ColumnType::Text => {
let end = bytes
.iter()
.position(|&b| b == 0)
.unwrap_or(VALUE_DATA_SIZE);
let s = std::str::from_utf8(&bytes[..end]).unwrap_or("");
Value::Text(crate::types::ArcString(std::sync::Arc::from(s)))
}
_ => Value::Null, }
}
fn value_to_bytes(&self, value: &Value) -> Result<[u8; VALUE_DATA_SIZE]> {
Self::value_to_bytes_helper(value)
}
fn value_to_bytes_helper(value: &Value) -> Result<[u8; VALUE_DATA_SIZE]> {
let mut buf = [0u8; VALUE_DATA_SIZE];
match value {
Value::Integer(i) => buf[..8].copy_from_slice(&i.to_be_bytes()),
Value::Float(f) => {
let canonical = if *f == 0.0 { 0.0f64 } else { *f }; let bits = canonical.to_bits();
let sortable = if bits & (1u64 << 63) != 0 {
!bits } else {
bits ^ (1u64 << 63) };
buf[..8].copy_from_slice(&sortable.to_be_bytes());
}
Value::Timestamp(ts) => buf[..8].copy_from_slice(&ts.as_micros().to_be_bytes()),
Value::Bool(b) => buf[0] = if *b { 1 } else { 0 },
Value::Text(s) => {
let raw = s.as_bytes();
let len = raw.len().min(VALUE_DATA_SIZE);
buf[..len].copy_from_slice(&raw[..len]);
}
_ => {
return Err(StorageError::InvalidData(format!(
"Unsupported value type for indexing: {:?}",
value
)));
}
};
Ok(buf)
}
}
#[derive(Debug, Clone)]
pub struct IndexStats {
pub cached_values: usize,
pub total_row_ids: usize,
}
use crate::index::builder::{BuildStats, IndexBuilder};
use crate::types::Row;
impl IndexBuilder for ColumnValueIndex {
fn build_from_memtable(&mut self, _rows: &[(RowId, Row)]) -> Result<()> {
debug_log!(
"[ColumnIndex::{}] ⚠️ build_from_memtable is deprecated, use insert_batch instead",
self.column_name
);
Ok(())
}
fn persist(&mut self) -> Result<()> {
use std::time::Instant;
let start = Instant::now();
self.flush()?;
let duration = start.elapsed();
debug_log!(
"[ColumnIndex::{}] Persist: {:?}",
self.column_name,
duration
);
Ok(())
}
fn name(&self) -> &str {
&self.column_name
}
fn stats(&self) -> BuildStats {
let stats = self.stats();
BuildStats {
rows_processed: stats.total_row_ids,
build_time_ms: 0,
persist_time_ms: 0,
index_size_bytes: stats.total_row_ids * IndexKey::key_size(),
}
}
}
impl ColumnValueIndex {
pub fn insert_batch(&self, batch: &[(RowId, &Value)]) -> Result<()> {
if batch.is_empty() {
return Ok(());
}
for (row_id, value) in batch {
self.insert(value, *row_id)?;
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::TempDir;
#[test]
fn test_column_value_index_basic() -> Result<()> {
let temp_dir = TempDir::new()?;
let path = temp_dir.path().join("test_index.idx");
let index = ColumnValueIndex::create(
&path,
"users".to_string(),
"age".to_string(),
ColumnValueIndexConfig::default(),
)?;
index.insert(&Value::Integer(25), 1)?;
index.insert(&Value::Integer(30), 2)?;
index.insert(&Value::Integer(25), 3)?;
let row_ids = index.get(&Value::Integer(25))?;
assert_eq!(row_ids.len(), 2);
assert!(row_ids.contains(&1));
assert!(row_ids.contains(&3));
Ok(())
}
#[test]
fn test_column_value_index_delete_tombstone() -> Result<()> {
let temp_dir = TempDir::new()?;
let path = temp_dir.path().join("test_tombstone.idx");
let index = ColumnValueIndex::create(
&path,
"users".to_string(),
"age".to_string(),
ColumnValueIndexConfig::default(),
)?;
index.insert(&Value::Integer(25), 1)?;
index.insert(&Value::Integer(25), 2)?;
index.delete(&Value::Integer(25), 1)?;
let row_ids = index.get(&Value::Integer(25))?;
assert_eq!(row_ids.len(), 1);
assert!(row_ids.contains(&2));
index.insert(&Value::Integer(25), 1)?;
let row_ids = index.get(&Value::Integer(25))?;
assert_eq!(row_ids.len(), 2);
assert!(row_ids.contains(&1));
assert!(row_ids.contains(&2));
Ok(())
}
#[test]
fn test_column_value_index_range_with_delete() -> Result<()> {
let temp_dir = TempDir::new()?;
let path = temp_dir.path().join("test_range_delete.idx");
let index = ColumnValueIndex::create(
&path,
"users".to_string(),
"age".to_string(),
ColumnValueIndexConfig::default(),
)?;
for i in 10..20 {
index.insert(&Value::Integer(i), i as RowId)?;
}
let deleted = index.delete_range(&Value::Integer(13), &Value::Integer(17))?;
assert!(deleted > 0);
let row_ids = index.range(&Value::Integer(10), &Value::Integer(19))?;
let expected: Vec<RowId> = vec![10, 11, 12, 18, 19];
assert_eq!(row_ids.len(), expected.len());
for id in &expected {
assert!(row_ids.contains(id));
}
Ok(())
}
#[test]
fn test_tombstone_key_normalization() {
let mut vb = [0u8; VALUE_DATA_SIZE];
vb[..5].copy_from_slice(b"hello");
let short = IndexKey {
value_bytes: vb,
row_id: 42,
};
let tk_short = tombstone_key(&short);
assert_eq!(tk_short.value_bytes, vb);
let mut vb2 = [0u8; VALUE_DATA_SIZE];
vb2[..12].copy_from_slice(b"abcdefghijkl");
let long = IndexKey {
value_bytes: vb2,
row_id: 99,
};
let tk_long = tombstone_key(&long);
assert_eq!(tk_long.value_bytes, vb2);
assert_eq!(tk_long.row_id, 99);
}
#[test]
fn test_column_value_index_long_text_tombstone() -> Result<()> {
let temp_dir = TempDir::new()?;
let path = temp_dir.path().join("test_long_text.idx");
let index = ColumnValueIndex::create(
&path,
"users".to_string(),
"bio".to_string(),
ColumnValueIndexConfig::default(),
)?;
let long_val = Value::text("abcdefghijklmno_xtralong_value".to_string());
index.insert(&long_val, 1)?;
index.insert(&long_val, 2)?;
index.flush_buffer()?;
index.delete(&long_val, 1)?;
let row_ids = index.get(&long_val)?;
assert_eq!(row_ids.len(), 1);
assert!(row_ids.contains(&2));
assert!(!row_ids.contains(&1));
Ok(())
}
#[test]
fn test_column_value_index_concurrent_stress() -> Result<()> {
use std::sync::atomic::{AtomicBool, Ordering};
let temp_dir = TempDir::new()?;
let path = temp_dir.path().join("test_concurrent.idx");
let index = Arc::new(ColumnValueIndex::create(
&path,
"users".to_string(),
"age".to_string(),
ColumnValueIndexConfig::default(),
)?);
let stop = Arc::new(AtomicBool::new(false));
let n = 500;
for i in 0..n {
index.insert(&Value::Integer(i % 50), i as RowId)?;
}
let mut handles = vec![];
{
let index = Arc::clone(&index);
let stop = Arc::clone(&stop);
handles.push(std::thread::spawn(move || {
while !stop.load(Ordering::Relaxed) {
for i in 0..100 {
let _ = index.insert(&Value::Integer(i % 50), i as RowId);
}
}
}));
}
{
let index = Arc::clone(&index);
let stop = Arc::clone(&stop);
handles.push(std::thread::spawn(move || {
while !stop.load(Ordering::Relaxed) {
for i in 0..50 {
let _ = index.delete(&Value::Integer(i), i as RowId);
let _ = index.insert(&Value::Integer(i), i as RowId);
}
}
}));
}
{
let index = Arc::clone(&index);
let stop = Arc::clone(&stop);
handles.push(std::thread::spawn(move || {
while !stop.load(Ordering::Relaxed) {
if let Ok(ids) = index.get(&Value::Integer(25)) {
for &id in &ids {
assert!(id < n as RowId, "get() returned unexpected row_id {}", id);
}
}
if let Ok(ids) = index.query_less_than_or_equal(&Value::Integer(10)) {
for &id in &ids {
assert!(id < n as RowId, "range() returned unexpected row_id {}", id);
}
}
}
}));
}
std::thread::sleep(std::time::Duration::from_millis(500));
stop.store(true, Ordering::Relaxed);
for handle in handles {
handle.join().unwrap();
}
for i in 0..10 {
index.delete(&Value::Integer(i), i as RowId)?;
}
for i in 0..10 {
let ids = index.get(&Value::Integer(i))?;
assert!(
!ids.contains(&(i as RowId)),
"Deleted key (value={}, row_id={}) still present",
i,
i
);
}
Ok(())
}
#[test]
fn test_concurrent_read_and_flush_no_deadlock() -> Result<()> {
use std::sync::Arc;
use std::time::Duration;
let temp_dir = TempDir::new()?;
let path = temp_dir.path().join("test_deadlock2.idx");
let index = Arc::new(ColumnValueIndex::create(
&path,
"t".to_string(),
"c".to_string(),
ColumnValueIndexConfig::default(),
)?);
for i in 0..2000i64 {
index.insert(&Value::Integer(i % 100), i as RowId)?;
}
let idx_reader = Arc::clone(&index);
let idx_writer = Arc::clone(&index);
let stop = Arc::new(std::sync::atomic::AtomicBool::new(false));
let s1 = Arc::clone(&stop);
let s2 = Arc::clone(&stop);
let reader = std::thread::spawn(move || {
while !s1.load(std::sync::atomic::Ordering::Relaxed) {
for v in 0..50i64 {
let _ = idx_reader.get_arc(&Value::Integer(v));
}
}
});
let writer = std::thread::spawn(move || {
while !s2.load(std::sync::atomic::Ordering::Relaxed) {
for i in 0..100i64 {
let _ = idx_writer.insert(&Value::Integer(i % 50), 10000 + i as RowId);
}
let _ = idx_writer.flush();
}
});
std::thread::sleep(Duration::from_secs(3));
stop.store(true, std::sync::atomic::Ordering::Relaxed);
reader.join().unwrap();
writer.join().unwrap();
eprintln!(" OK: concurrent read+flush deadlock regression passed");
Ok(())
}
#[test]
fn test_update_moves_row_id() -> Result<()> {
let temp_dir = TempDir::new()?;
let path = temp_dir.path().join("test_update.idx");
let index = ColumnValueIndex::create(
&path,
"t".to_string(),
"c".to_string(),
ColumnValueIndexConfig::default(),
)?;
index.insert(&Value::Integer(100), 1)?;
assert_eq!(index.get(&Value::Integer(100))?, vec![1]);
assert!(index.get(&Value::Integer(200))?.is_empty());
index.update(&Value::Integer(100), &Value::Integer(200), 1)?;
assert!(
!index.get(&Value::Integer(100))?.contains(&1),
"old value should not contain updated row"
);
assert!(
index.get(&Value::Integer(200))?.contains(&1),
"new value should contain updated row"
);
assert_eq!(
index.get(&Value::Integer(100))?.len() + index.get(&Value::Integer(200))?.len(),
1
);
Ok(())
}
#[test]
fn test_update_same_value_noop() -> Result<()> {
let temp_dir = TempDir::new()?;
let path = temp_dir.path().join("test_update_same.idx");
let index = ColumnValueIndex::create(
&path,
"t".to_string(),
"c".to_string(),
ColumnValueIndexConfig::default(),
)?;
index.insert(&Value::Integer(100), 5)?;
index.update(&Value::Integer(100), &Value::Integer(100), 5)?;
assert!(
index.get(&Value::Integer(100))?.contains(&5),
"row should still be present after same-value update"
);
Ok(())
}
#[test]
fn test_insert_flush_get() -> Result<()> {
let temp_dir = TempDir::new()?;
let path = temp_dir.path().join("test_flush_get.idx");
let index = ColumnValueIndex::create(
&path,
"t".to_string(),
"c".to_string(),
ColumnValueIndexConfig::default(),
)?;
for i in 0..500i64 {
index.insert(&Value::Integer(i % 20), i as RowId)?;
}
index.flush()?;
let ids = index.get(&Value::Integer(5))?;
assert!(!ids.is_empty(), "data should survive flush");
let range_ids = index.range(&Value::Integer(0), &Value::Integer(10))?;
assert!(!range_ids.is_empty(), "range query should work after flush");
Ok(())
}
#[test]
fn test_batch_insert_and_flush() -> Result<()> {
let temp_dir = TempDir::new()?;
let path = temp_dir.path().join("test_batch2.idx");
let index = ColumnValueIndex::create(
&path,
"t".to_string(),
"c".to_string(),
ColumnValueIndexConfig::default(),
)?;
let items: Vec<(Value, RowId)> = (0..1000i64)
.map(|i| (Value::Integer(i % 10), i as RowId))
.collect();
index.batch_insert(items)?;
index.flush()?;
for v in 0..10i64 {
assert_eq!(
index.get(&Value::Integer(v))?.len(),
100,
"value {} should have 100 row_ids",
v
);
}
Ok(())
}
#[test]
fn test_delete_makes_entry_invisible() -> Result<()> {
let temp_dir = TempDir::new()?;
let path = temp_dir.path().join("test_del2.idx");
let index = ColumnValueIndex::create(
&path,
"t".to_string(),
"c".to_string(),
ColumnValueIndexConfig::default(),
)?;
index.insert(&Value::Integer(42), 100)?;
index.delete(&Value::Integer(42), 100)?;
assert!(index.get(&Value::Integer(42))?.is_empty());
Ok(())
}
#[test]
fn test_update_survives_flush() -> Result<()> {
let temp_dir = TempDir::new()?;
let path = temp_dir.path().join("test_upd_flush2.idx");
let index = ColumnValueIndex::create(
&path,
"t".to_string(),
"c".to_string(),
ColumnValueIndexConfig::default(),
)?;
index.insert(&Value::Integer(10), 1)?;
index.update(&Value::Integer(10), &Value::Integer(20), 1)?;
index.flush()?;
assert!(
index.get(&Value::Integer(20))?.contains(&1),
"after update+flush: row 1 should be at new value"
);
assert!(
!index.get(&Value::Integer(10))?.contains(&1),
"after update+flush: row 1 should NOT be at old value"
);
Ok(())
}
#[test]
fn test_query_between_exclusive_boundaries() -> Result<()> {
let temp_dir = TempDir::new()?;
let path = temp_dir.path().join("test_between.idx");
let index = ColumnValueIndex::create(
path,
"t".to_string(),
"v".to_string(),
ColumnValueIndexConfig::default(),
)?;
index.insert(&Value::Integer(10), 0)?;
index.insert(&Value::Integer(20), 1)?;
index.insert(&Value::Integer(30), 2)?;
let result = index.query_between(&Value::Integer(10), true, &Value::Integer(30), true)?;
assert_eq!(
result.len(),
3,
"[10,30] inclusive should find 3, got {}",
result.len()
);
let result = index.query_between(&Value::Integer(10), false, &Value::Integer(30), false)?;
assert_eq!(
result.len(),
1,
"(10,30) exclusive should find 1, got {}",
result.len()
);
assert!(result.contains(&1), "should contain row_id=1 (value=20)");
let result = index.query_between(&Value::Integer(10), false, &Value::Integer(30), true)?;
assert_eq!(
result.len(),
2,
"(10,30] should find 2, got {}",
result.len()
);
let result = index.query_between(&Value::Integer(10), true, &Value::Integer(30), false)?;
assert_eq!(
result.len(),
2,
"[10,30) should find 2, got {}",
result.len()
);
Ok(())
}
}