use super::*;
use crate::error::{DbError, DbResult};
use crate::storage::index::{extract_field_value, generate_ngrams, tokenize};
use crate::storage::serializer::deserialize_doc;
use rust_rocksdb::{Direction, IteratorMode, WriteBatch};
use serde_json::Value;
use hex;
type IndexEntry = (Vec<u8>, Vec<u8>);
type IndexEntries = Vec<IndexEntry>;
type IndexKeys = Vec<Vec<u8>>;
fn is_already_exists(e: &DbError) -> bool {
matches!(e, DbError::InvalidDocument(msg) if msg.contains("already exists"))
}
fn is_not_found(e: &DbError) -> bool {
matches!(e, DbError::InvalidDocument(msg) if msg.contains("not found"))
}
impl Collection {
pub fn create_index_from_spec(&self, spec: &crate::storage::index::IndexSpec) -> DbResult<()> {
use crate::storage::index::IndexSpec;
match spec {
IndexSpec::Regular {
name,
fields,
index_type,
unique,
} => self
.create_index(name.clone(), fields.clone(), index_type.clone(), *unique)
.map(|_| ()),
IndexSpec::Fulltext {
name,
fields,
min_length,
} => self
.create_fulltext_index(name.clone(), fields.clone(), *min_length)
.map(|_| ()),
IndexSpec::Geo { name, field } => self
.create_geo_index(name.clone(), field.clone())
.map(|_| ()),
IndexSpec::Ttl {
name,
field,
expire_after_seconds,
} => self
.create_ttl_index(name.clone(), field.clone(), *expire_after_seconds)
.map(|_| ()),
IndexSpec::Vector(config) => self.create_vector_index(config.clone()).map(|_| ()),
}
}
pub fn apply_index_spec(&self, spec: &crate::storage::index::IndexSpec) -> DbResult<()> {
match self.create_index_from_spec(spec) {
Ok(()) => Ok(()),
Err(e) if is_already_exists(&e) => Ok(()),
Err(e) => Err(e),
}
}
pub fn drop_index_of_kind(
&self,
kind: crate::storage::index::IndexKind,
name: &str,
) -> DbResult<()> {
use crate::storage::index::IndexKind;
match kind {
IndexKind::Regular => self.drop_index(name),
IndexKind::Geo => self.drop_geo_index(name),
IndexKind::Ttl => self.drop_ttl_index(name),
IndexKind::Vector => self.drop_vector_index(name),
}
}
pub fn apply_index_drop(
&self,
kind: crate::storage::index::IndexKind,
name: &str,
) -> DbResult<()> {
match self.drop_index_of_kind(kind, name) {
Ok(()) => Ok(()),
Err(e) if is_not_found(&e) => Ok(()),
Err(e) => Err(e),
}
}
pub fn get_all_indexes(&self) -> Vec<Index> {
self.index_meta()
.expect("Column family should exist")
.indexes
.clone()
}
pub(crate) fn get_index(&self, name: &str) -> Option<Index> {
self.index_meta()?
.indexes
.iter()
.find(|i| i.name == name)
.cloned()
}
pub fn create_index(
&self,
name: String,
fields: Vec<String>,
index_type: IndexType,
unique: bool,
) -> DbResult<IndexStats> {
if self.get_index(&name).is_some() {
return Err(DbError::InvalidDocument(format!(
"Index '{}' already exists",
name
)));
}
let index = Index::new(name.clone(), fields.clone(), index_type.clone(), unique);
let index_bytes = serde_json::to_vec(&index)?;
{
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.expect("Column family should exist");
db.put_cf(&cf, Self::idx_meta_key(&name), &index_bytes)
.map_err(|e| DbError::InternalError(format!("Failed to create index: {}", e)))?;
}
self.invalidate_index_meta();
let indexed_count = match self.backfill_index(&name, &fields, &index_type) {
Ok(count) => count,
Err(e) => {
if let Err(cleanup) = self.drop_index(&name) {
tracing::error!(
collection = %self.name,
index = %name,
error = %cleanup,
"index backfill failed and the partial index could not be dropped; \
drop it by hand before querying the field"
);
}
return Err(e);
}
};
if index_type == IndexType::Bloom {
if let Some(filter) = self.bloom_filters.get(&name) {
self.save_bloom_filter(&name, &filter)?;
}
} else if index_type == IndexType::Cuckoo {
if let Some(filter) = self.cuckoo_filters.get(&name) {
self.save_cuckoo_filter(&name, &filter)?;
}
}
Ok(IndexStats {
name,
field: fields.first().cloned().unwrap_or_default(),
fields,
index_type,
unique,
unique_values: indexed_count, indexed_documents: indexed_count,
})
}
fn backfill_index(
&self,
name: &str,
fields: &[String],
index_type: &IndexType,
) -> DbResult<usize> {
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.ok_or_else(|| DbError::CollectionNotFound(self.name.clone()))?;
let prefix = DOC_PREFIX.as_bytes();
let mut batch = WriteBatch::default();
let mut indexed_count = 0;
const BATCH_SIZE_LIMIT: usize = 10000;
for item in db.prefix_iterator_cf(&cf, prefix) {
let (key, value) = item.map_err(|e| {
DbError::InternalError(format!("Failed to read documents for index: {}", e))
})?;
if !key.starts_with(prefix) {
break;
}
let Ok(doc) = deserialize_doc(&value) else {
continue;
};
let doc_value = doc.to_value();
let field_values: Vec<Value> = fields
.iter()
.map(|f| extract_field_value(&doc_value, f))
.collect();
if field_values.iter().all(|v| v.is_null()) {
continue;
}
let entry_key = Self::idx_entry_key(name, &field_values, &doc.key);
batch.put_cf(&cf, entry_key, doc.key.as_bytes());
indexed_count += 1;
if *index_type == IndexType::Bloom {
for value in &field_values {
self.bloom_insert(name, &value.to_string());
}
} else if *index_type == IndexType::Cuckoo {
for value in &field_values {
self.cuckoo_insert(name, &value.to_string());
}
}
if indexed_count % BATCH_SIZE_LIMIT == 0 {
db.write(&batch).map_err(|e| {
DbError::InternalError(format!("Failed to write index batch: {}", e))
})?;
batch = WriteBatch::default();
}
}
if !batch.is_empty() {
db.write(&batch).map_err(|e| {
DbError::InternalError(format!("Failed to write final index batch: {}", e))
})?;
}
Ok(indexed_count)
}
pub fn drop_index(&self, name: &str) -> DbResult<()> {
if self.get_index(name).is_none() {
return Err(DbError::InvalidDocument(format!(
"Index '{}' not found",
name
)));
}
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.expect("Column family should exist");
let mut batch = WriteBatch::default();
let mut deleted_count = 0;
const BATCH_SIZE_LIMIT: usize = 10000;
batch.delete_cf(&cf, Self::idx_meta_key(name));
let prefix = format!("{}{}:", IDX_PREFIX, name);
let iter = db.prefix_iterator_cf(&cf, prefix.as_bytes());
for result in iter.flatten() {
let (key, _) = result;
if key.starts_with(prefix.as_bytes()) {
batch.delete_cf(&cf, &key);
deleted_count += 1;
if deleted_count % BATCH_SIZE_LIMIT == 0 {
db.write(&batch).map_err(|e| {
DbError::InternalError(format!("Failed to write drop batch: {}", e))
})?;
batch = WriteBatch::default();
}
} else {
break;
}
}
if deleted_count > 0 || !batch.is_empty() {
db.write(&batch).map_err(|e| {
DbError::InternalError(format!("Failed to write final drop batch: {}", e))
})?;
}
self.invalidate_index_meta();
Ok(())
}
pub fn list_indexes(&self) -> Vec<IndexStats> {
let mut stats: Vec<IndexStats> = self
.get_all_indexes()
.iter()
.filter_map(|idx| self.get_index_stats(&idx.name))
.collect();
for idx in self.get_all_fulltext_indexes() {
stats.push(IndexStats {
name: idx.name,
fields: idx.fields.clone(),
field: idx.fields.first().cloned().unwrap_or_default(),
index_type: IndexType::Fulltext,
unique: false,
unique_values: 0, indexed_documents: 0, });
}
stats
}
pub fn rebuild_all_indexes(&self) -> DbResult<usize> {
let total_start = std::time::Instant::now();
let indexes = self.get_all_indexes();
let geo_indexes = self.get_all_geo_indexes();
let ft_indexes = self.get_all_fulltext_indexes();
tracing::info!(
"rebuild_all_indexes: {} regular, {} geo, {} fulltext indexes",
indexes.len(),
geo_indexes.len(),
ft_indexes.len()
);
if indexes.is_empty() && geo_indexes.is_empty() && ft_indexes.is_empty() {
return Ok(0);
}
let clear_start = std::time::Instant::now();
{
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.expect("Column family should exist");
let mut batch = WriteBatch::default();
for index in &indexes {
let prefix = format!("{}{}:", IDX_PREFIX, index.name);
let iter = db.prefix_iterator_cf(&cf, prefix.as_bytes());
for result in iter.flatten() {
let (key, _) = result;
if key.starts_with(prefix.as_bytes()) {
batch.delete_cf(&cf, &key);
} else {
break;
}
}
}
for geo_index in &geo_indexes {
let prefix = format!("{}{}:", GEO_PREFIX, geo_index.name);
let iter = db.prefix_iterator_cf(&cf, prefix.as_bytes());
for result in iter.flatten() {
let (key, _) = result;
if key.starts_with(prefix.as_bytes()) {
batch.delete_cf(&cf, &key);
} else {
break;
}
}
}
for ft_index in &ft_indexes {
let ngram_prefix = format!("{}{}:", FT_PREFIX, ft_index.name);
let iter = db.prefix_iterator_cf(&cf, ngram_prefix.as_bytes());
for result in iter.flatten() {
let (key, _) = result;
if key.starts_with(ngram_prefix.as_bytes()) {
batch.delete_cf(&cf, &key);
} else {
break;
}
}
let term_prefix = format!("{}{}:", FT_TERM_PREFIX, ft_index.name);
let iter = db.prefix_iterator_cf(&cf, term_prefix.as_bytes());
for result in iter.flatten() {
let (key, _) = result;
if key.starts_with(term_prefix.as_bytes()) {
batch.delete_cf(&cf, &key);
} else {
break;
}
}
}
let _ = db.write(&batch);
}
tracing::info!(
"rebuild_all_indexes: Clear phase took {:?}",
clear_start.elapsed()
);
let stream_start = std::time::Instant::now();
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.expect("Column family should exist");
let prefix = DOC_PREFIX.as_bytes();
let iter = db.prefix_iterator_cf(&cf, prefix);
const BATCH_SIZE_LIMIT: usize = 10000; let mut doc_count = 0;
let mut regular_batch = if !indexes.is_empty() {
Some(WriteBatch::default())
} else {
None
};
let mut geo_batch = if !geo_indexes.is_empty() {
Some(WriteBatch::default())
} else {
None
};
let mut ft_batch = if !ft_indexes.is_empty() {
Some(WriteBatch::default())
} else {
None
};
let mut regular_count = 0;
let mut geo_count = 0;
let mut ft_count = 0;
for result in iter.flatten() {
let (key, value) = result;
if !key.starts_with(prefix) {
break;
}
if let Ok(doc) = deserialize_doc(&value) {
doc_count += 1;
let doc_value = doc.to_value();
if let Some(ref mut batch) = regular_batch {
for index in &indexes {
let field_values: Vec<Value> = index
.fields
.iter()
.map(|f| extract_field_value(&doc_value, f))
.collect();
if !field_values.iter().all(|v| v.is_null()) {
let entry_key =
Self::idx_entry_key(&index.name, &field_values, &doc.key);
batch.put_cf(&cf, entry_key, doc.key.as_bytes());
regular_count += 1;
if regular_count % BATCH_SIZE_LIMIT == 0 {
let _ = db.write(&*batch);
*batch = WriteBatch::default();
}
}
}
}
if let Some(ref mut batch) = geo_batch {
for geo_index in &geo_indexes {
let field_value = extract_field_value(&doc_value, &geo_index.field);
if !field_value.is_null() {
let entry_key = Self::geo_entry_key(&geo_index.name, &doc.key);
if let Ok(geo_data) = serde_json::to_vec(&field_value) {
batch.put_cf(&cf, entry_key, &geo_data);
geo_count += 1;
if geo_count % BATCH_SIZE_LIMIT == 0 {
let _ = db.write(&*batch);
*batch = WriteBatch::default();
}
}
}
}
}
if let Some(ref mut batch) = ft_batch {
for ft_index in &ft_indexes {
for field in &ft_index.fields {
let field_value = extract_field_value(&doc_value, field);
if let Some(text) = field_value.as_str() {
let terms = tokenize(text);
for term in &terms {
if term.len() >= ft_index.min_length {
let term_key =
Self::ft_term_key(&ft_index.name, term, &doc.key);
batch.put_cf(&cf, term_key, doc.key.as_bytes());
ft_count += 1;
}
}
let ngrams = generate_ngrams(text, NGRAM_SIZE);
for ngram in &ngrams {
let ngram_key =
Self::ft_ngram_key(&ft_index.name, ngram, &doc.key);
batch.put_cf(&cf, ngram_key, doc.key.as_bytes());
ft_count += 1;
}
if ft_count >= BATCH_SIZE_LIMIT {
let _ = db.write(&*batch);
*batch = WriteBatch::default();
ft_count = 0;
}
}
}
}
}
}
}
if let Some(batch) = regular_batch {
if regular_count > 0 && regular_count % BATCH_SIZE_LIMIT != 0 {
let _ = db.write(&batch);
}
}
if let Some(batch) = geo_batch {
if geo_count > 0 && geo_count % BATCH_SIZE_LIMIT != 0 {
let _ = db.write(&batch);
}
}
if let Some(batch) = ft_batch {
if ft_count > 0 {
let _ = db.write(&batch);
}
}
tracing::info!(
"rebuild_all_indexes: Streamed {} docs in {:?}",
doc_count,
stream_start.elapsed()
);
let vector_configs = self.get_all_vector_index_configs();
if !vector_configs.is_empty() {
let vec_start = std::time::Instant::now();
for entry in self.vector_indexes.iter() {
entry.clear();
}
let prefix = DOC_PREFIX.as_bytes();
let iter = db.prefix_iterator_cf(&cf, prefix);
for result in iter.flatten() {
let (key, value) = result;
if !key.starts_with(prefix) {
break;
}
if let Ok(doc) = deserialize_doc(&value) {
let doc_value = doc.to_value();
self.update_vector_indexes_on_upsert(&doc.key, &doc_value);
}
}
if let Err(e) = self.persist_vector_indexes() {
tracing::warn!("Failed to persist vector indexes during rebuild: {}", e);
}
tracing::info!(
"rebuild_all_indexes: Vector indexes ({}) took {:?}",
vector_configs.len(),
vec_start.elapsed()
);
}
tracing::info!(
"rebuild_all_indexes: Total time {:?}",
total_start.elapsed()
);
Ok(doc_count)
}
pub fn index_documents(&self, docs: &[Document]) -> DbResult<usize> {
let total_start = std::time::Instant::now();
let indexes = self.get_all_indexes();
let geo_indexes = self.get_all_geo_indexes();
let ft_indexes = self.get_all_fulltext_indexes();
if indexes.is_empty() && geo_indexes.is_empty() && ft_indexes.is_empty() {
return Ok(0);
}
for index in &indexes {
if index.index_type == IndexType::Bloom {
let _ = self.get_or_create_bloom_filter(&index.name);
} else if index.index_type == IndexType::Cuckoo {
self.preload_cuckoo_filter(&index.name);
}
}
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.expect("Column family should exist");
if !indexes.is_empty() {
let mut batch = WriteBatch::default();
for doc in docs {
let doc_value = doc.to_value();
for index in &indexes {
let field_values: Vec<Value> = index
.fields
.iter()
.map(|f| extract_field_value(&doc_value, f))
.collect();
if !field_values.iter().all(|v| v.is_null()) {
let entry_key = Self::idx_entry_key(&index.name, &field_values, &doc.key);
batch.put_cf(&cf, entry_key, doc.key.as_bytes());
if index.index_type == IndexType::Bloom {
for value in &field_values {
self.bloom_insert(&index.name, &value.to_string());
}
} else if index.index_type == IndexType::Cuckoo {
for value in &field_values {
self.cuckoo_insert(&index.name, &value.to_string());
}
}
}
}
}
let _ = db.write(&batch);
}
if !geo_indexes.is_empty() {
let mut batch = WriteBatch::default();
for doc in docs {
let doc_value = doc.to_value();
for geo_index in &geo_indexes {
let field_value = extract_field_value(&doc_value, &geo_index.field);
if !field_value.is_null() {
let entry_key = Self::geo_entry_key(&geo_index.name, &doc.key);
if let Ok(geo_data) = serde_json::to_vec(&field_value) {
batch.put_cf(&cf, entry_key, &geo_data);
}
}
}
}
let _ = db.write(&batch);
}
if !ft_indexes.is_empty() {
let mut batch = WriteBatch::default();
for doc in docs {
let doc_value = doc.to_value();
for ft_index in &ft_indexes {
for field in &ft_index.fields {
let field_value = extract_field_value(&doc_value, field);
if let Some(text) = field_value.as_str() {
let terms = tokenize(text);
for term in &terms {
if term.len() >= ft_index.min_length {
let term_key =
Self::ft_term_key(&ft_index.name, term, &doc.key);
batch.put_cf(&cf, term_key, doc.key.as_bytes());
}
}
let ngrams = generate_ngrams(text, NGRAM_SIZE);
for ngram in &ngrams {
let ngram_key = Self::ft_ngram_key(&ft_index.name, ngram, &doc.key);
batch.put_cf(&cf, ngram_key, doc.key.as_bytes());
}
}
}
}
}
let _ = db.write(&batch);
}
tracing::info!("index_documents: Total time {:?}", total_start.elapsed());
Ok(docs.len())
}
pub fn get_index_stats(&self, name: &str) -> Option<IndexStats> {
let index = self.get_index(name)?;
let db = &self.db;
let cf = db.cf_handle(&self.name)?;
let prefix = format!("{}{}:", IDX_PREFIX, name);
let iter = db.prefix_iterator_cf(&cf, prefix.as_bytes());
let count = iter
.filter(|r| {
r.as_ref()
.map(|(k, _)| k.starts_with(prefix.as_bytes()))
.unwrap_or(false)
})
.count();
Some(IndexStats {
name: index.name,
fields: index.fields.clone(),
field: index.fields.first().cloned().unwrap_or_default(),
index_type: index.index_type,
unique: index.unique,
unique_values: count,
indexed_documents: count,
})
}
pub(crate) fn check_unique_constraints(
&self,
doc_key: &str,
doc_value: &Value,
) -> DbResult<()> {
let indexes = self.get_all_indexes();
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.expect("Column family should exist");
for index in indexes {
if index.unique {
let field_values: Vec<Value> = index
.fields
.iter()
.map(|f| extract_field_value(doc_value, f))
.collect();
if field_values.iter().all(|v| v.is_null()) {
continue;
}
let encoded_values: Vec<String> = field_values
.iter()
.map(|v| hex::encode(crate::storage::codec::encode_key(v)))
.collect();
let value_part = encoded_values.join("_");
let prefix = format!("{}{}:{}:", IDX_PREFIX, index.name, value_part);
let mut iter = db.prefix_iterator_cf(&cf, prefix.as_bytes());
if let Some(Ok((key, value))) = iter.next() {
if key.starts_with(prefix.as_bytes()) {
let existing_key = String::from_utf8_lossy(&value); if existing_key != doc_key {
return Err(DbError::InvalidDocument(format!(
"Unique constraint violated: fields '{:?}' with value {:?} already exists in index '{}'",
index.fields, field_values, index.name
)));
}
}
}
}
}
Ok(())
}
#[allow(dead_code)]
pub(crate) fn update_indexes_on_insert(
&self,
doc_key: &str,
doc_value: &Value,
) -> DbResult<()> {
let indexes = self.get_all_indexes();
for index in &indexes {
if index.index_type == IndexType::Bloom {
let _ = self.get_or_create_bloom_filter(&index.name);
} else if index.index_type == IndexType::Cuckoo {
self.preload_cuckoo_filter(&index.name);
}
}
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.expect("Column family should exist");
for index in indexes {
let field_values: Vec<Value> = index
.fields
.iter()
.map(|f| extract_field_value(doc_value, f))
.collect();
if !field_values.iter().all(|v| v.is_null()) {
let entry_key = Self::idx_entry_key(&index.name, &field_values, doc_key);
db.put_cf(&cf, entry_key, doc_key.as_bytes()).map_err(|e| {
DbError::InternalError(format!("Failed to update index: {}", e))
})?;
if index.index_type == IndexType::Bloom {
for value in field_values {
self.bloom_insert(&index.name, &value.to_string());
}
} else if index.index_type == IndexType::Cuckoo {
for value in field_values {
self.cuckoo_insert(&index.name, &value.to_string());
}
}
}
}
let _ = db;
let geo_indexes = self.get_all_geo_indexes();
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.expect("Column family should exist");
for geo_index in geo_indexes {
let field_value = extract_field_value(doc_value, &geo_index.field);
if !field_value.is_null() {
let entry_key = Self::geo_entry_key(&geo_index.name, doc_key);
let geo_data = serde_json::to_vec(&field_value)?;
db.put_cf(&cf, entry_key, &geo_data).map_err(|e| {
DbError::InternalError(format!("Failed to update geo index: {}", e))
})?;
}
}
Ok(())
}
#[allow(dead_code)]
pub(crate) fn update_indexes_on_update(
&self,
doc_key: &str,
old_value: &Value,
new_value: &Value,
) -> DbResult<()> {
let indexes = self.get_all_indexes();
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.expect("Column family should exist");
for index in indexes {
let old_values: Vec<Value> = index
.fields
.iter()
.map(|f| extract_field_value(old_value, f))
.collect();
let new_values: Vec<Value> = index
.fields
.iter()
.map(|f| extract_field_value(new_value, f))
.collect();
if !old_values.iter().all(|v| v.is_null()) {
let old_entry_key = Self::idx_entry_key(&index.name, &old_values, doc_key);
db.delete_cf(&cf, old_entry_key).map_err(|e| {
DbError::InternalError(format!("Failed to update index: {}", e))
})?;
}
if !new_values.iter().all(|v| v.is_null()) {
let new_entry_key = Self::idx_entry_key(&index.name, &new_values, doc_key);
db.put_cf(&cf, new_entry_key, doc_key.as_bytes())
.map_err(|e| {
DbError::InternalError(format!("Failed to update index: {}", e))
})?;
}
}
let _ = db;
let geo_indexes = self.get_all_geo_indexes();
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.expect("Column family should exist");
for geo_index in geo_indexes {
let entry_key = Self::geo_entry_key(&geo_index.name, doc_key);
let new_field = extract_field_value(new_value, &geo_index.field);
if !new_field.is_null() {
let geo_data = serde_json::to_vec(&new_field)?;
db.put_cf(&cf, entry_key, &geo_data).map_err(|e| {
DbError::InternalError(format!("Failed to update geo index: {}", e))
})?;
} else {
db.delete_cf(&cf, entry_key).map_err(|e| {
DbError::InternalError(format!("Failed to update geo index: {}", e))
})?;
}
}
Ok(())
}
#[allow(dead_code)]
pub(crate) fn update_indexes_on_delete(
&self,
doc_key: &str,
doc_value: &Value,
) -> DbResult<()> {
let indexes = self.get_all_indexes();
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.expect("Column family should exist");
for index in indexes {
let field_values: Vec<Value> = index
.fields
.iter()
.map(|f| extract_field_value(doc_value, f))
.collect();
if !field_values.iter().all(|v| v.is_null()) {
let entry_key = Self::idx_entry_key(&index.name, &field_values, doc_key);
db.delete_cf(&cf, entry_key).map_err(|e| {
DbError::InternalError(format!("Failed to update index: {}", e))
})?;
}
}
let _ = db;
let geo_indexes = self.get_all_geo_indexes();
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.expect("Column family should exist");
for geo_index in geo_indexes {
let entry_key = Self::geo_entry_key(&geo_index.name, doc_key);
db.delete_cf(&cf, entry_key).map_err(|e| {
DbError::InternalError(format!("Failed to update geo index: {}", e))
})?;
}
Ok(())
}
pub fn get_index_for_field(&self, field: &str) -> Option<Index> {
self.index_meta()?
.indexes
.iter()
.find(|index| index.fields.first().map(|s| s.as_str()) == Some(field))
.cloned()
}
fn get_single_field_index_for_field(&self, field: &str) -> Option<Index> {
self.index_meta()?
.indexes
.iter()
.find(|index| {
index.fields.len() == 1 && index.fields.first().map(|s| s.as_str()) == Some(field)
})
.cloned()
}
pub fn index_lookup_gt(
&self,
field: &str,
value: &Value,
limit: Option<usize>,
) -> Option<Vec<Document>> {
self.index_range_scan(field, value, false, true, limit)
}
pub fn index_lookup_gte(
&self,
field: &str,
value: &Value,
limit: Option<usize>,
) -> Option<Vec<Document>> {
self.index_range_scan(field, value, true, true, limit)
}
pub fn index_lookup_lt(
&self,
field: &str,
value: &Value,
limit: Option<usize>,
) -> Option<Vec<Document>> {
self.index_range_scan(field, value, false, false, limit)
}
pub fn index_lookup_lte(
&self,
field: &str,
value: &Value,
limit: Option<usize>,
) -> Option<Vec<Document>> {
self.index_range_scan(field, value, true, false, limit)
}
fn index_range_scan(
&self,
field: &str,
value: &Value,
inclusive: bool,
forward: bool,
limit: Option<usize>,
) -> Option<Vec<Document>> {
let index = self.get_single_field_index_for_field(field)?;
let index_name = &index.name;
let value_key = crate::storage::codec::encode_key(value);
let value_str = hex::encode(value_key);
let db = &self.db;
let cf = db.cf_handle(&self.name)?;
let prefix_base = format!("{}{}:", IDX_PREFIX, index_name);
let seek_key = format!("{}{}:", prefix_base, value_str);
let cap = limit.unwrap_or(1000);
let mut doc_keys = Vec::new();
if forward {
let mode = IteratorMode::From(seek_key.as_bytes(), Direction::Forward);
let iter = db.iterator_cf(&cf, mode);
for result in iter {
if let Ok((k, v)) = result {
if !k.starts_with(prefix_base.as_bytes()) {
break;
}
if !inclusive && k.starts_with(seek_key.as_bytes()) {
continue;
}
let val_str = String::from_utf8_lossy(&v);
doc_keys.push(Self::doc_key(&val_str));
}
if doc_keys.len() >= cap {
break;
}
}
} else {
let mode = IteratorMode::From(seek_key.as_bytes(), Direction::Reverse);
let iter = db.iterator_cf(&cf, mode);
for result in iter {
if let Ok((k, v)) = result {
if !k.starts_with(prefix_base.as_bytes()) {
break;
}
if !inclusive && k.starts_with(seek_key.as_bytes()) {
continue;
}
let val_str = String::from_utf8_lossy(&v);
doc_keys.push(Self::doc_key(&val_str));
}
if doc_keys.len() >= cap {
break;
}
}
}
if doc_keys.is_empty() {
return Some(Vec::new());
}
let results = db.multi_get_cf(doc_keys.iter().map(|k| (&cf, k.as_slice())));
let docs: Vec<Document> = results
.into_iter()
.filter_map(|r| r.ok())
.flatten()
.filter_map(|bytes| deserialize_doc(&bytes).ok())
.collect();
Some(docs)
}
pub fn index_lookup_eq(&self, field: &str, value: &Value) -> Option<Vec<Document>> {
if value.is_null() {
return None;
}
let index = self.get_single_field_index_for_field(field)?;
if (index.index_type == IndexType::Bloom
&& !self.bloom_check(&index.name, &value.to_string()))
|| (index.index_type == IndexType::Cuckoo
&& !self.cuckoo_check(&index.name, &value.to_string()))
{
return Some(Vec::new());
}
let value_str = hex::encode(crate::storage::codec::encode_key(value));
let db = &self.db;
let cf = db.cf_handle(&self.name)?;
let prefix = format!("{}{}:{}:", IDX_PREFIX, index.name, value_str);
let iter = db.prefix_iterator_cf(&cf, prefix.as_bytes());
let doc_keys: Vec<Vec<u8>> = iter
.filter_map(|r| r.ok())
.take_while(|(k, _)| k.starts_with(prefix.as_bytes()))
.map(|(_, v)| {
let key_str = String::from_utf8_lossy(&v);
Self::doc_key(&key_str)
})
.collect();
if doc_keys.is_empty() {
return Some(Vec::new());
}
let results = db.multi_get_cf(doc_keys.iter().map(|k| (&cf, k.as_slice())));
let docs: Vec<Document> = results
.into_iter()
.filter_map(|r| r.ok())
.flatten()
.filter_map(|bytes| deserialize_doc(&bytes).ok())
.collect();
Some(docs)
}
pub fn index_lookup_eq_composite(
&self,
conditions: &[(String, Value)],
) -> Option<(crate::storage::index::Index, Vec<Document>)> {
if conditions.len() < 2 {
return None;
}
if conditions.iter().any(|(_, v)| v.is_null()) {
return None;
}
let mut best: Option<crate::storage::index::Index> = None;
for index in self.get_all_indexes() {
if index.fields.len() < 2 {
continue;
}
let all_covered = index
.fields
.iter()
.all(|f| conditions.iter().any(|(cf, _)| cf == f));
if all_covered
&& best
.as_ref()
.is_none_or(|b| b.fields.len() < index.fields.len())
{
best = Some(index);
}
}
let index = best?;
let encoded: Vec<String> = index
.fields
.iter()
.map(|f| {
let v = conditions
.iter()
.find(|(cf, _)| cf == f)
.map(|(_, v)| v)
.expect("coverage check ensures field is present");
hex::encode(crate::storage::codec::encode_key(v))
})
.collect();
let value_part = encoded.join("_");
let db = &self.db;
let cf = db.cf_handle(&self.name)?;
let prefix = format!("{}{}:{}:", IDX_PREFIX, index.name, value_part);
let iter = db.prefix_iterator_cf(&cf, prefix.as_bytes());
let doc_keys: Vec<Vec<u8>> = iter
.filter_map(|r| r.ok())
.take_while(|(k, _)| k.starts_with(prefix.as_bytes()))
.map(|(_, v)| {
let key_str = String::from_utf8_lossy(&v);
Self::doc_key(&key_str)
})
.collect();
if doc_keys.is_empty() {
return Some((index, Vec::new()));
}
let results = db.multi_get_cf(doc_keys.iter().map(|k| (&cf, k.as_slice())));
let docs: Vec<Document> = results
.into_iter()
.filter_map(|r| r.ok())
.flatten()
.filter_map(|bytes| deserialize_doc(&bytes).ok())
.collect();
Some((index, docs))
}
pub fn index_lookup_eq_limit(
&self,
field: &str,
value: &Value,
limit: usize,
) -> Option<Vec<Document>> {
if value.is_null() {
return None;
}
let index = self.get_single_field_index_for_field(field)?;
if (index.index_type == IndexType::Bloom
&& !self.bloom_check(&index.name, &value.to_string()))
|| (index.index_type == IndexType::Cuckoo
&& !self.cuckoo_check(&index.name, &value.to_string()))
{
return Some(Vec::new());
}
let value_str = hex::encode(crate::storage::codec::encode_key(value));
let db = &self.db;
let cf = db.cf_handle(&self.name)?;
let prefix = format!("{}{}:{}:", IDX_PREFIX, index.name, value_str);
let iter = db.prefix_iterator_cf(&cf, prefix.as_bytes());
let doc_keys: Vec<Vec<u8>> = iter
.filter_map(|r| r.ok())
.take_while(|(k, _)| k.starts_with(prefix.as_bytes()))
.take(limit)
.map(|(_, v)| {
let key_str = String::from_utf8_lossy(&v);
Self::doc_key(&key_str)
})
.collect();
if doc_keys.is_empty() {
return Some(Vec::new());
}
let results = db.multi_get_cf(doc_keys.iter().map(|k| (&cf, k.as_slice())));
let docs: Vec<Document> = results
.into_iter()
.filter_map(|r| r.ok())
.flatten()
.filter_map(|bytes| deserialize_doc(&bytes).ok())
.collect();
Some(docs)
}
pub fn index_sorted(
&self,
field: &str,
ascending: bool,
limit: Option<usize>,
) -> Option<Vec<Document>> {
if field == "_id" || field == "_key" {
let db = &self.db;
let cf = db.cf_handle(&self.name)?;
let prefix = DOC_PREFIX.as_bytes();
let iter = if ascending {
let mode = IteratorMode::From(prefix, Direction::Forward);
db.iterator_cf(&cf, mode)
} else {
let mut seek_key = prefix.to_vec();
seek_key.push(0xFF);
let mode = IteratorMode::From(&seek_key, Direction::Reverse);
db.iterator_cf(&cf, mode)
};
let docs: Vec<Document> = iter
.filter_map(|r| r.ok())
.take_while(|(k, _)| k.starts_with(prefix))
.filter_map(|(_, v)| deserialize_doc(&v).ok())
.take(limit.unwrap_or(usize::MAX))
.collect();
return Some(docs);
}
let index = self.get_index_for_field(field)?;
let index_name = index.name.clone();
let prefix = format!("{}{}:", IDX_PREFIX, index_name);
let db = &self.db;
let cf = db.cf_handle(&self.name)?;
let prefix_bytes = prefix.as_bytes();
let iter = if ascending {
let mode = IteratorMode::From(prefix_bytes, Direction::Forward);
db.iterator_cf(&cf, mode)
} else {
let mut seek_key = prefix.as_bytes().to_vec();
seek_key.push(0xFF);
let mode = IteratorMode::From(&seek_key, Direction::Reverse);
db.iterator_cf(&cf, mode)
};
let doc_keys: Vec<String> = iter
.filter_map(|r| r.ok())
.take_while(|(k, _)| k.starts_with(prefix_bytes))
.map(|(_, v)| String::from_utf8_lossy(&v).to_string())
.take(limit.unwrap_or(usize::MAX))
.collect();
let _ = db;
if doc_keys.is_empty() {
return Some(Vec::new());
}
let docs = self.get_many(&doc_keys);
let doc_map: std::collections::HashMap<_, _> =
docs.into_iter().map(|d| (d.key.clone(), d)).collect();
let result: Vec<Document> = doc_keys
.into_iter()
.filter_map(|key| doc_map.get(&key).cloned())
.collect();
Some(result)
}
pub(crate) fn get_or_create_bloom_filter(&self, index_name: &str) -> DbResult<BloomFilter> {
if let Some(filter) = self.bloom_filters.get(index_name) {
return Ok(filter.clone());
}
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.ok_or(DbError::InternalError("CF not found".into()))?;
let key = format!("{}{}", BLO_IDX_PREFIX, index_name);
let new_filter = if let Ok(Some(bytes)) = db.get_cf(&cf, key.as_bytes()) {
if let Ok(filter) = serde_json::from_slice::<BloomFilter>(&bytes) {
filter
} else {
BloomFilter::with_num_bits(1024 * 8).expected_items(1000)
}
} else {
BloomFilter::with_num_bits(1024 * 8).expected_items(1000)
};
self.bloom_filters
.insert(index_name.to_string(), new_filter.clone());
Ok(new_filter)
}
pub(crate) fn save_bloom_filter(&self, index_name: &str, filter: &BloomFilter) -> DbResult<()> {
let db = &self.db;
let cf = db
.cf_handle(&self.name)
.ok_or(DbError::InternalError("CF not found".into()))?;
let key = format!("{}{}", BLO_IDX_PREFIX, index_name);
let bytes = serde_json::to_vec(filter)?;
db.put_cf(&cf, key.as_bytes(), &bytes)
.map_err(|e| DbError::InternalError(e.to_string()))?;
Ok(())
}
pub(crate) fn bloom_insert(&self, index_name: &str, item: &str) {
if let Ok(filter) = self.get_or_create_bloom_filter(index_name) {
let mut filter = filter;
filter.insert(item.as_bytes());
self.bloom_filters.insert(index_name.to_string(), filter);
}
}
pub fn bloom_check(&self, index_name: &str, item: &str) -> bool {
if let Ok(filter) = self.get_or_create_bloom_filter(index_name) {
filter.contains(item.as_bytes())
} else {
true }
}
pub(crate) fn preload_cuckoo_filter(&self, index_name: &str) {
if self.cuckoo_filters.contains_key(index_name) {
return;
}
let db = &self.db;
if let Some(cf) = db.cf_handle(&self.name) {
let key = format!("{}{}", CFO_IDX_PREFIX, index_name);
if let Ok(Some(_bytes)) = db.get_cf(&cf, key.as_bytes()) {
self.cuckoo_filters
.insert(index_name.to_string(), CuckooFilter::new());
} else {
self.cuckoo_filters
.insert(index_name.to_string(), CuckooFilter::new());
}
}
}
pub(crate) fn save_cuckoo_filter(
&self,
index_name: &str,
_filter: &CuckooFilter<DefaultHasher>,
) -> DbResult<()> {
let db = &self.db;
let _cf = db
.cf_handle(&self.name)
.ok_or(DbError::InternalError("CF not found".into()))?;
let _key = format!("{}{}", CFO_IDX_PREFIX, index_name);
Ok(())
}
pub fn cuckoo_insert(&self, index_name: &str, item: &str) {
self.preload_cuckoo_filter(index_name);
if let Some(mut filter) = self.cuckoo_filters.get_mut(index_name) {
let _ = filter.add(item);
}
}
pub fn cuckoo_delete(&self, index_name: &str, item: &str) {
self.preload_cuckoo_filter(index_name);
if let Some(mut filter) = self.cuckoo_filters.get_mut(index_name) {
let _ = filter.delete(item);
}
}
pub fn cuckoo_check(&self, index_name: &str, item: &str) -> bool {
self.preload_cuckoo_filter(index_name);
if let Some(filter) = self.cuckoo_filters.get(index_name) {
filter.contains(item)
} else {
true
}
}
pub(crate) fn compute_index_entries_for_insert(
&self,
doc_key: &str,
doc_value: &Value,
) -> DbResult<(IndexEntries, IndexEntries)> {
let indexes = self.get_all_indexes();
let mut regular_entries = Vec::new();
for index in indexes {
let field_values: Vec<Value> = index
.fields
.iter()
.map(|f| extract_field_value(doc_value, f))
.collect();
if !field_values.iter().all(|v| v.is_null()) {
let entry_key = Self::idx_entry_key(&index.name, &field_values, doc_key);
regular_entries.push((entry_key, doc_key.as_bytes().to_vec()));
}
}
let geo_indexes = self.get_all_geo_indexes();
let mut geo_entries = Vec::new();
for geo_index in &geo_indexes {
let field_value = extract_field_value(doc_value, &geo_index.field);
if !field_value.is_null() {
let entry_key = Self::geo_entry_key(&geo_index.name, doc_key);
let geo_data = serde_json::to_vec(&field_value)?;
geo_entries.push((entry_key, geo_data));
}
}
Ok((regular_entries, geo_entries))
}
pub(crate) fn compute_index_entries_for_update(
&self,
doc_key: &str,
old_value: &Value,
new_value: &Value,
) -> DbResult<(IndexEntries, IndexKeys, IndexEntries, IndexKeys)> {
let indexes = self.get_all_indexes();
let mut entries_to_add = Vec::new();
let mut keys_to_remove = Vec::new();
for index in indexes {
let old_values: Vec<Value> = index
.fields
.iter()
.map(|f| extract_field_value(old_value, f))
.collect();
let new_values: Vec<Value> = index
.fields
.iter()
.map(|f| extract_field_value(new_value, f))
.collect();
if !old_values.iter().all(|v| v.is_null()) {
let old_entry_key = Self::idx_entry_key(&index.name, &old_values, doc_key);
keys_to_remove.push(old_entry_key);
}
if !new_values.iter().all(|v| v.is_null()) {
let new_entry_key = Self::idx_entry_key(&index.name, &new_values, doc_key);
entries_to_add.push((new_entry_key, doc_key.as_bytes().to_vec()));
}
}
let geo_indexes = self.get_all_geo_indexes();
let mut geo_entries_to_add = Vec::new();
let mut geo_keys_to_remove = Vec::new();
for geo_index in &geo_indexes {
let entry_key = Self::geo_entry_key(&geo_index.name, doc_key);
let new_field = extract_field_value(new_value, &geo_index.field);
geo_keys_to_remove.push(entry_key.clone());
if !new_field.is_null() {
let geo_data = serde_json::to_vec(&new_field)?;
geo_entries_to_add.push((entry_key, geo_data));
}
}
Ok((
entries_to_add,
keys_to_remove,
geo_entries_to_add,
geo_keys_to_remove,
))
}
pub(crate) fn compute_index_entries_for_delete(
&self,
doc_key: &str,
doc_value: &Value,
) -> DbResult<(IndexKeys, IndexKeys)> {
let indexes = self.get_all_indexes();
let mut keys_to_remove = Vec::new();
for index in indexes {
let field_values: Vec<Value> = index
.fields
.iter()
.map(|f| extract_field_value(doc_value, f))
.collect();
if !field_values.iter().all(|v| v.is_null()) {
let entry_key = Self::idx_entry_key(&index.name, &field_values, doc_key);
keys_to_remove.push(entry_key);
}
}
let geo_indexes = self.get_all_geo_indexes();
let mut geo_keys_to_remove = Vec::new();
for geo_index in &geo_indexes {
let entry_key = Self::geo_entry_key(&geo_index.name, doc_key);
geo_keys_to_remove.push(entry_key);
}
Ok((keys_to_remove, geo_keys_to_remove))
}
}