use parking_lot::RwLock;
use std::collections::BTreeMap;
use std::ops::Bound;
use std::sync::atomic::{AtomicBool, Ordering as AtomicOrdering};
use rustc_hash::FxHashMap;
use crate::expression::Expression;
use crate::traits::Index;
use radixdb_core::{CompactArc, CompactVec, I64Map};
use radixdb_core::{DataType, Error, IndexEntry, IndexType, Operator, Result, RowIdVec, Value};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CompositeKey(pub Vec<Value>);
impl PartialOrd for CompositeKey {
fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
Some(self.cmp(other))
}
}
impl Ord for CompositeKey {
fn cmp(&self, other: &Self) -> std::cmp::Ordering {
for (a, b) in self.0.iter().zip(other.0.iter()) {
match a.cmp(b) {
std::cmp::Ordering::Equal => continue,
ord => return ord,
}
}
self.0.len().cmp(&other.0.len())
}
}
impl std::hash::Hash for CompositeKey {
fn hash<H: std::hash::Hasher>(&self, state: &mut H) {
for v in &self.0 {
v.hash(state);
}
}
}
pub struct MultiColumnIndex {
name: String,
table_name: String,
column_names: Vec<String>,
column_ids: Vec<i32>,
data_types: Vec<DataType>,
is_unique: bool,
closed: AtomicBool,
publication_fence: RwLock<()>,
sorted_values: RwLock<BTreeMap<CompositeKey, CompactVec<i64>>>,
btree_built: AtomicBool,
value_to_rows: RwLock<FxHashMap<CompositeKey, CompactVec<i64>>>,
prefix_indexes: Vec<RwLock<FxHashMap<CompositeKey, CompactVec<i64>>>>,
prefix_built: Vec<AtomicBool>,
row_to_key: RwLock<I64Map<Vec<CompactArc<Value>>>>,
}
impl std::fmt::Debug for MultiColumnIndex {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("MultiColumnIndex")
.field("name", &self.name)
.field("table_name", &self.table_name)
.field("column_names", &self.column_names)
.field("column_ids", &self.column_ids)
.field("is_unique", &self.is_unique)
.field(
"btree_built",
&self.btree_built.load(AtomicOrdering::Relaxed),
)
.field("closed", &self.closed.load(AtomicOrdering::Relaxed))
.finish_non_exhaustive()
}
}
impl MultiColumnIndex {
#[doc(hidden)]
pub fn new(
name: String,
table_name: String,
column_names: Vec<String>,
column_ids: Vec<i32>,
data_types: Vec<DataType>,
is_unique: bool,
expected_rows: usize,
) -> Self {
Self::try_new(
name,
table_name,
column_names,
column_ids,
data_types,
is_unique,
expected_rows,
)
.expect("internal MultiColumnIndex metadata must be coherent")
}
pub fn try_new(
name: String,
table_name: String,
column_names: Vec<String>,
column_ids: Vec<i32>,
data_types: Vec<DataType>,
is_unique: bool,
expected_rows: usize,
) -> Result<Self> {
super::validate_index_metadata_shape(
&name,
&table_name,
&column_names,
&column_ids,
&data_types,
)?;
let num_cols = column_names.len();
let mut prefix_indexes = Vec::with_capacity(num_cols);
let mut prefix_built = Vec::with_capacity(num_cols);
for _ in 0..num_cols {
prefix_indexes.push(RwLock::new(FxHashMap::default()));
prefix_built.push(AtomicBool::new(false));
}
Ok(Self {
name,
table_name,
column_names,
column_ids,
data_types,
is_unique,
closed: AtomicBool::new(false),
publication_fence: RwLock::new(()),
sorted_values: RwLock::new(BTreeMap::new()),
btree_built: AtomicBool::new(false),
value_to_rows: RwLock::new(if expected_rows > 0 {
FxHashMap::with_capacity_and_hasher(expected_rows, Default::default())
} else {
FxHashMap::default()
}),
prefix_indexes,
prefix_built,
row_to_key: RwLock::new(if expected_rows > 0 {
I64Map::with_capacity(expected_rows)
} else {
I64Map::new()
}),
})
}
#[inline]
fn values_match(stored: &[CompactArc<Value>], input: &[Value]) -> bool {
if stored.len() != input.len() {
return false;
}
match stored.len() {
0 => true,
1 => stored[0].as_ref() == &input[0],
2 => stored[0].as_ref() == &input[0] && stored[1].as_ref() == &input[1],
3 => {
stored[0].as_ref() == &input[0]
&& stored[1].as_ref() == &input[1]
&& stored[2].as_ref() == &input[2]
}
4 => {
stored[0].as_ref() == &input[0]
&& stored[1].as_ref() == &input[1]
&& stored[2].as_ref() == &input[2]
&& stored[3].as_ref() == &input[3]
}
_ => stored
.iter()
.zip(input.iter())
.all(|(s, i)| s.as_ref() == i),
}
}
#[inline]
fn range_is_empty(
min: &[Value],
max: &[Value],
min_inclusive: bool,
max_inclusive: bool,
) -> bool {
if min.is_empty() || max.is_empty() || min.len() != max.len() {
return false;
}
match min.cmp(max) {
std::cmp::Ordering::Greater => true,
std::cmp::Ordering::Equal => !min_inclusive || !max_inclusive,
std::cmp::Ordering::Less => false,
}
}
fn ensure_btree_built(&self) {
if self.btree_built.load(AtomicOrdering::Acquire) {
return;
}
let value_to_rows = self.value_to_rows.read();
let mut sorted_values = self.sorted_values.write();
if self.btree_built.load(AtomicOrdering::Acquire) {
return;
}
for (key, rows) in value_to_rows.iter() {
sorted_values.insert(key.clone(), rows.clone());
}
self.btree_built.store(true, AtomicOrdering::Release);
}
fn ensure_prefix_built(&self, prefix_len: usize) {
if prefix_len == 0 || prefix_len > self.prefix_built.len() {
return;
}
let idx = prefix_len - 1;
if self.prefix_built[idx].load(AtomicOrdering::Acquire) {
return;
}
let row_to_key = self.row_to_key.read();
let mut prefix_index = self.prefix_indexes[idx].write();
if self.prefix_built[idx].load(AtomicOrdering::Acquire) {
return;
}
for (row_id, arc_values) in row_to_key.iter() {
if arc_values.len() >= prefix_len {
let prefix_key = CompositeKey(
arc_values[..prefix_len]
.iter()
.map(|a| (**a).clone())
.collect(),
);
let rows = prefix_index.entry(prefix_key).or_default();
if let Err(pos) = rows.binary_search(&row_id) {
rows.insert(pos, row_id);
}
}
}
self.prefix_built[idx].store(true, AtomicOrdering::Release);
}
fn check_unique_constraint_locked(
&self,
key: &CompositeKey,
row_id: i64,
value_to_rows: &FxHashMap<CompositeKey, CompactVec<i64>>,
) -> Result<()> {
if !self.is_unique {
return Ok(());
}
for v in &key.0 {
if v.is_null() {
return Ok(());
}
}
if let Some(rows) = value_to_rows.get(key) {
if !rows.is_empty() && !rows.contains(&row_id) {
let values_str: Vec<String> = key.0.iter().map(|v| format!("{:?}", v)).collect();
return Err(Error::unique_constraint(
&self.name,
self.column_names.join(", "),
format!("[{}]", values_str.join(", ")),
));
}
}
Ok(())
}
}
impl Index for MultiColumnIndex {
fn name(&self) -> &str {
&self.name
}
fn table_name(&self) -> &str {
&self.table_name
}
fn build(&mut self) -> Result<()> {
Ok(())
}
fn add(&self, values: &[Value], row_id: i64, _ref_id: i64) -> Result<()> {
let _publication = self.publication_fence.write();
if self.closed.load(AtomicOrdering::Acquire) {
return Err(Error::IndexClosed);
}
let num_cols = self.column_ids.len();
if values.len() != num_cols {
return Err(Error::internal(format!(
"expected {} values, got {}",
num_cols,
values.len()
)));
}
let key = CompositeKey(values.to_vec());
let mut old_key_for_cleanup: Option<Vec<CompactArc<Value>>> = None;
let mut value_to_rows = self.value_to_rows.write();
let mut row_to_key = self.row_to_key.write();
if let Some(existing_arc_values) = row_to_key.get(row_id) {
if Self::values_match(existing_arc_values, values) {
return Ok(());
}
old_key_for_cleanup = Some(existing_arc_values.clone());
}
self.check_unique_constraint_locked(&key, row_id, &value_to_rows)?;
if let Some(ref old_arc_values) = old_key_for_cleanup {
let existing_key = CompositeKey(old_arc_values.iter().map(|a| (**a).clone()).collect());
if let Some(rows) = value_to_rows.get_mut(&existing_key) {
rows.retain(|id| *id != row_id);
if rows.is_empty() {
value_to_rows.remove(&existing_key);
}
}
}
let btree_needs_update = self.btree_built.load(AtomicOrdering::Acquire);
let arc_values: Vec<CompactArc<Value>> =
values.iter().map(|v| CompactArc::new(v.clone())).collect();
let key_for_hash = key.clone();
let hash_rows = value_to_rows.entry(key_for_hash).or_default();
if let Err(pos) = hash_rows.binary_search(&row_id) {
hash_rows.insert(pos, row_id);
}
row_to_key.insert(row_id, arc_values);
drop(value_to_rows);
drop(row_to_key);
if btree_needs_update {
let mut sorted_values = self.sorted_values.write();
if let Some(ref old_arc_values) = old_key_for_cleanup {
let existing_key =
CompositeKey(old_arc_values.iter().map(|a| (**a).clone()).collect());
if let Some(rows) = sorted_values.get_mut(&existing_key) {
rows.retain(|id| *id != row_id);
if rows.is_empty() {
sorted_values.remove(&existing_key);
}
}
}
let btree_rows = sorted_values.entry(key).or_default();
if let Err(pos) = btree_rows.binary_search(&row_id) {
btree_rows.insert(pos, row_id);
}
}
for prefix_len in 1..num_cols {
let idx = prefix_len - 1;
if self.prefix_built[idx].load(AtomicOrdering::Acquire) {
let mut prefix_index = self.prefix_indexes[idx].write();
if let Some(ref old_arc_values) = old_key_for_cleanup {
if old_arc_values.len() >= prefix_len {
let old_prefix_key = CompositeKey(
old_arc_values[..prefix_len]
.iter()
.map(|a| (**a).clone())
.collect(),
);
if let Some(rows) = prefix_index.get_mut(&old_prefix_key) {
rows.retain(|id| *id != row_id);
if rows.is_empty() {
prefix_index.remove(&old_prefix_key);
}
}
}
}
let prefix_key = CompositeKey(values[..prefix_len].to_vec());
let prefix_rows = prefix_index.entry(prefix_key).or_default();
if let Err(pos) = prefix_rows.binary_search(&row_id) {
prefix_rows.insert(pos, row_id);
}
}
}
Ok(())
}
fn add_batch(&self, entries: &I64Map<Vec<Value>>) -> Result<()> {
let borrowed: Vec<(i64, &[Value])> = entries
.iter()
.map(|(row_id, values)| (row_id, values.as_slice()))
.collect();
self.add_batch_slice(&borrowed)
}
fn remove(&self, values: &[Value], row_id: i64, _ref_id: i64) -> Result<()> {
let _publication = self.publication_fence.write();
if self.closed.load(AtomicOrdering::Acquire) {
return Err(Error::IndexClosed);
}
let stored_values = {
let mut value_to_rows = self.value_to_rows.write();
let mut row_to_key = self.row_to_key.write();
let Some(stored) = row_to_key.get(row_id).cloned() else {
return Ok(());
};
if !Self::values_match(&stored, values) {
return Err(Error::invalid_argument(format!(
"stale values supplied while removing row {row_id} from index '{}'",
self.name
)));
}
let stored_values: Vec<Value> = stored.iter().map(|value| (**value).clone()).collect();
let key = CompositeKey(stored_values.clone());
if let Some(rows) = value_to_rows.get_mut(&key) {
if let Ok(pos) = rows.binary_search(&row_id) {
rows.remove(pos);
}
if rows.is_empty() {
value_to_rows.remove(&key);
}
}
row_to_key.remove(row_id);
stored_values
};
let key = CompositeKey(stored_values.clone());
if self.btree_built.load(AtomicOrdering::Acquire) {
let mut sorted_values = self.sorted_values.write();
if let Some(rows) = sorted_values.get_mut(&key) {
if let Ok(pos) = rows.binary_search(&row_id) {
rows.remove(pos);
}
if rows.is_empty() {
sorted_values.remove(&key);
}
}
}
for prefix_len in 1..stored_values.len() {
let idx = prefix_len - 1;
if self.prefix_built[idx].load(AtomicOrdering::Acquire) {
let prefix_key = CompositeKey(stored_values[..prefix_len].to_vec());
let mut prefix_index = self.prefix_indexes[idx].write();
if let Some(rows) = prefix_index.get_mut(&prefix_key) {
if let Ok(pos) = rows.binary_search(&row_id) {
rows.remove(pos);
}
if rows.is_empty() {
prefix_index.remove(&prefix_key);
}
}
}
}
Ok(())
}
fn remove_batch(&self, entries: &I64Map<Vec<Value>>) -> Result<()> {
let borrowed: Vec<(i64, &[Value])> = entries
.iter()
.map(|(row_id, values)| (row_id, values.as_slice()))
.collect();
self.remove_batch_slice(&borrowed)
}
fn add_batch_slice(&self, entries: &[(i64, &[Value])]) -> Result<()> {
let _publication = self.publication_fence.write();
if entries.is_empty() {
return Ok(());
}
if self.closed.load(AtomicOrdering::Acquire) {
return Err(Error::IndexClosed);
}
let num_cols = self.column_ids.len();
for &(_row_id, values) in entries {
if values.len() != num_cols {
return Err(Error::internal(format!(
"expected {} values, got {}",
num_cols,
values.len()
)));
}
}
let mut value_to_rows = self.value_to_rows.write();
let mut row_to_key = self.row_to_key.write();
value_to_rows.reserve(entries.len());
row_to_key.reserve(entries.len());
let btree_needs_update = self.btree_built.load(AtomicOrdering::Acquire);
if self.is_unique {
let mut batch_keys: FxHashMap<CompositeKey, i64> =
FxHashMap::with_capacity_and_hasher(entries.len(), Default::default());
for &(row_id, values) in entries {
let has_null = values.iter().any(|v| v.is_null());
if has_null {
continue;
}
let key = CompositeKey(values.to_vec());
if let Some(&existing_row_id) = batch_keys.get(&key) {
if existing_row_id != row_id {
let values_str: Vec<String> =
values.iter().map(|v| format!("{:?}", v)).collect();
return Err(Error::unique_constraint(
&self.name,
self.column_names.join(", "),
format!("[{}]", values_str.join(", ")),
));
}
}
if let Some(existing_rows) = value_to_rows.get(&key) {
if !existing_rows.is_empty() && !existing_rows.contains(&row_id) {
let values_str: Vec<String> =
values.iter().map(|v| format!("{:?}", v)).collect();
return Err(Error::unique_constraint(
&self.name,
self.column_names.join(", "),
format!("[{}]", values_str.join(", ")),
));
}
}
batch_keys.insert(key, row_id);
}
}
let mut updates_to_old_key: Vec<(i64, Vec<CompactArc<Value>>)> = Vec::new();
for &(row_id, values) in entries {
let key = CompositeKey(values.to_vec());
let arc_values: Vec<CompactArc<Value>> =
values.iter().map(|v| CompactArc::new(v.clone())).collect();
if let Some(existing_arc_values) = row_to_key.get(row_id) {
if !Self::values_match(existing_arc_values, values) {
updates_to_old_key.push((row_id, existing_arc_values.clone()));
let existing_key =
CompositeKey(existing_arc_values.iter().map(|a| (**a).clone()).collect());
if let Some(rows) = value_to_rows.get_mut(&existing_key) {
rows.retain(|id| *id != row_id);
if rows.is_empty() {
value_to_rows.remove(&existing_key);
}
}
}
}
let hash_rows = value_to_rows.entry(key).or_default();
if let Err(pos) = hash_rows.binary_search(&row_id) {
hash_rows.insert(pos, row_id);
}
row_to_key.insert(row_id, arc_values);
}
drop(value_to_rows);
drop(row_to_key);
if btree_needs_update {
let mut sorted_values = self.sorted_values.write();
for (row_id, existing_arc_values) in &updates_to_old_key {
let existing_key =
CompositeKey(existing_arc_values.iter().map(|a| (**a).clone()).collect());
if let Some(rows) = sorted_values.get_mut(&existing_key) {
rows.retain(|id| *id != *row_id);
if rows.is_empty() {
sorted_values.remove(&existing_key);
}
}
}
for &(row_id, values) in entries {
let key = CompositeKey(values.to_vec());
let btree_rows = sorted_values.entry(key).or_default();
if let Err(pos) = btree_rows.binary_search(&row_id) {
btree_rows.insert(pos, row_id);
}
}
}
for prefix_len in 1..num_cols {
let idx = prefix_len - 1;
if self.prefix_built[idx].load(AtomicOrdering::Acquire) {
let mut prefix_index = self.prefix_indexes[idx].write();
for (row_id, existing_arc_values) in &updates_to_old_key {
if existing_arc_values.len() >= prefix_len {
let prefix_key = CompositeKey(
existing_arc_values[..prefix_len]
.iter()
.map(|a| (**a).clone())
.collect(),
);
if let Some(rows) = prefix_index.get_mut(&prefix_key) {
rows.retain(|id| *id != *row_id);
if rows.is_empty() {
prefix_index.remove(&prefix_key);
}
}
}
}
for &(row_id, values) in entries {
let prefix_key = CompositeKey(values[..prefix_len].to_vec());
let prefix_rows = prefix_index.entry(prefix_key).or_default();
if let Err(pos) = prefix_rows.binary_search(&row_id) {
prefix_rows.insert(pos, row_id);
}
}
}
}
Ok(())
}
fn remove_batch_slice(&self, entries: &[(i64, &[Value])]) -> Result<()> {
let _publication = self.publication_fence.write();
if entries.is_empty() {
return Ok(());
}
if self.closed.load(AtomicOrdering::Acquire) {
return Err(Error::IndexClosed);
}
let mut value_to_rows = self.value_to_rows.write();
let mut row_to_key = self.row_to_key.write();
let mut removals = Vec::with_capacity(entries.len());
for &(row_id, values) in entries {
let Some(stored) = row_to_key.get(row_id) else {
continue;
};
if !Self::values_match(stored, values) {
return Err(Error::invalid_argument(format!(
"stale values supplied while removing row {row_id} from index '{}'",
self.name
)));
}
removals.push((
row_id,
stored
.iter()
.map(|value| (**value).clone())
.collect::<Vec<_>>(),
));
}
for (row_id, values) in &removals {
let key = CompositeKey(values.clone());
if let Some(rows) = value_to_rows.get_mut(&key) {
if let Ok(pos) = rows.binary_search(row_id) {
rows.remove(pos);
}
if rows.is_empty() {
value_to_rows.remove(&key);
}
}
row_to_key.remove(*row_id);
}
drop(value_to_rows);
drop(row_to_key);
if self.btree_built.load(AtomicOrdering::Acquire) {
let mut sorted_values = self.sorted_values.write();
for (row_id, values) in &removals {
let key = CompositeKey(values.clone());
if let Some(rows) = sorted_values.get_mut(&key) {
if let Ok(pos) = rows.binary_search(row_id) {
rows.remove(pos);
}
if rows.is_empty() {
sorted_values.remove(&key);
}
}
}
}
for prefix_len in 1..self.column_ids.len() {
let idx = prefix_len - 1;
if self.prefix_built[idx].load(AtomicOrdering::Acquire) {
let mut prefix_index = self.prefix_indexes[idx].write();
for (row_id, values) in &removals {
if values.len() >= prefix_len {
let prefix_key = CompositeKey(values[..prefix_len].to_vec());
if let Some(rows) = prefix_index.get_mut(&prefix_key) {
if let Ok(pos) = rows.binary_search(row_id) {
rows.remove(pos);
}
if rows.is_empty() {
prefix_index.remove(&prefix_key);
}
}
}
}
}
}
Ok(())
}
fn column_ids(&self) -> &[i32] {
&self.column_ids
}
fn column_names(&self) -> &[String] {
&self.column_names
}
fn data_types(&self) -> &[DataType] {
&self.data_types
}
fn index_type(&self) -> IndexType {
IndexType::MultiColumn
}
fn is_unique(&self) -> bool {
self.is_unique
}
fn find(&self, values: &[Value]) -> Result<Vec<IndexEntry>> {
let _publication = self.publication_fence.read();
if self.closed.load(AtomicOrdering::Acquire) {
return Err(Error::IndexClosed);
}
if values.is_empty() || values.len() > self.column_ids.len() {
return Err(Error::internal("invalid value count"));
}
let key = CompositeKey(values.to_vec());
if values.len() == self.column_ids.len() {
let value_to_rows = self.value_to_rows.read();
if let Some(row_ids) = value_to_rows.get(&key) {
return Ok(row_ids
.iter()
.map(|&row_id| IndexEntry { row_id, ref_id: 0 })
.collect());
}
return Ok(vec![]);
}
self.ensure_prefix_built(values.len());
let prefix_index = self.prefix_indexes[values.len() - 1].read();
if let Some(row_ids) = prefix_index.get(&key) {
return Ok(row_ids
.iter()
.map(|&row_id| IndexEntry { row_id, ref_id: 0 })
.collect());
}
Ok(vec![])
}
fn find_range(
&self,
min: &[Value],
max: &[Value],
min_inclusive: bool,
max_inclusive: bool,
) -> Result<Vec<IndexEntry>> {
let _publication = self.publication_fence.read();
if self.closed.load(AtomicOrdering::Acquire) {
return Err(Error::IndexClosed);
}
if min.len() > self.column_ids.len() || max.len() > self.column_ids.len() {
return Err(Error::invalid_argument(
"composite range bound exceeds index column count",
));
}
if Self::range_is_empty(min, max, min_inclusive, max_inclusive) {
return Ok(Vec::new());
}
self.ensure_btree_built();
let sorted_values = self.sorted_values.read();
let mut results = Vec::new();
let min_key = CompositeKey(min.to_vec());
let max_key = CompositeKey(max.to_vec());
let start = if min.is_empty() {
Bound::Unbounded
} else if min_inclusive {
Bound::Included(min_key.clone())
} else {
Bound::Excluded(min_key.clone())
};
let max_is_partial = !max.is_empty() && max.len() < self.column_ids.len();
let end = if max.is_empty() || max_is_partial {
Bound::Unbounded
} else if max_inclusive {
Bound::Included(max_key.clone())
} else {
Bound::Excluded(max_key.clone())
};
for (key, row_ids) in sorted_values.range((start, end)) {
let mut matches = true;
if !min.is_empty() && min.len() < key.0.len() {
let cmp = key.0[..min.len()].cmp(min);
if cmp == std::cmp::Ordering::Less
|| (!min_inclusive && cmp == std::cmp::Ordering::Equal)
{
matches = false;
}
}
if matches && !max.is_empty() && max.len() < key.0.len() {
let cmp = key.0[..max.len()].cmp(max);
if cmp == std::cmp::Ordering::Greater {
break;
}
if !max_inclusive && cmp == std::cmp::Ordering::Equal {
matches = false;
}
}
if matches {
for &row_id in row_ids {
results.push(IndexEntry { row_id, ref_id: 0 });
}
}
}
Ok(results)
}
fn find_range_ordered_limited(
&self,
min: &[Value],
max: &[Value],
min_inclusive: bool,
max_inclusive: bool,
ascending: bool,
limit: usize,
) -> Result<Vec<IndexEntry>> {
let _publication = self.publication_fence.read();
if self.closed.load(AtomicOrdering::Acquire) {
return Err(Error::IndexClosed);
}
if min.len() > self.column_ids.len() || max.len() > self.column_ids.len() {
return Err(Error::invalid_argument(
"composite range bound exceeds index column count",
));
}
if limit == 0 {
return Ok(Vec::new());
}
if Self::range_is_empty(min, max, min_inclusive, max_inclusive) {
return Ok(Vec::new());
}
self.ensure_btree_built();
let sorted_values = self.sorted_values.read();
let min_key = CompositeKey(min.to_vec());
let max_key = CompositeKey(max.to_vec());
let start = if min.is_empty() {
Bound::Unbounded
} else if min_inclusive {
Bound::Included(min_key)
} else {
Bound::Excluded(min_key)
};
let max_is_partial = !max.is_empty() && max.len() < self.column_ids.len();
let end = if max.is_empty() || max_is_partial {
Bound::Unbounded
} else if max_inclusive {
Bound::Included(max_key)
} else {
Bound::Excluded(max_key)
};
let mut results = Vec::with_capacity(limit.min(1024));
if ascending {
for (key, row_ids) in sorted_values.range((start, end)) {
if !min.is_empty() && min.len() < key.0.len() {
let cmp = key.0[..min.len()].cmp(min);
if cmp == std::cmp::Ordering::Less
|| (!min_inclusive && cmp == std::cmp::Ordering::Equal)
{
continue;
}
}
if !max.is_empty() && max.len() < key.0.len() {
let cmp = key.0[..max.len()].cmp(max);
if cmp == std::cmp::Ordering::Greater {
break;
}
if !max_inclusive && cmp == std::cmp::Ordering::Equal {
continue;
}
}
for &row_id in row_ids {
results.push(IndexEntry { row_id, ref_id: 0 });
if results.len() == limit {
return Ok(results);
}
}
}
} else {
for (key, row_ids) in sorted_values.range((start, end)).rev() {
if !max.is_empty() && max.len() < key.0.len() {
let cmp = key.0[..max.len()].cmp(max);
if cmp == std::cmp::Ordering::Greater
|| (!max_inclusive && cmp == std::cmp::Ordering::Equal)
{
continue;
}
}
if !min.is_empty() && min.len() < key.0.len() {
let cmp = key.0[..min.len()].cmp(min);
if cmp == std::cmp::Ordering::Less {
break;
}
if !min_inclusive && cmp == std::cmp::Ordering::Equal {
continue;
}
}
for &row_id in row_ids.iter().rev() {
results.push(IndexEntry { row_id, ref_id: 0 });
if results.len() == limit {
return Ok(results);
}
}
}
}
Ok(results)
}
fn find_with_operator(&self, op: Operator, values: &[Value]) -> Result<Vec<IndexEntry>> {
match op {
Operator::Eq => self.find(values),
Operator::Lt => self.find_range(&[], values, false, false),
Operator::Lte => self.find_range(&[], values, false, true),
Operator::Gt => self.find_range(values, &[], false, false),
Operator::Gte => self.find_range(values, &[], true, false),
_ => Err(Error::internal(format!("unsupported operator {:?}", op))),
}
}
fn get_row_ids_equal_into(&self, values: &[Value], buffer: &mut Vec<i64>) -> Result<()> {
let _publication = self.publication_fence.read();
if self.closed.load(AtomicOrdering::Acquire) {
return Err(Error::IndexClosed);
}
if values.is_empty() || values.len() > self.column_ids.len() {
return Err(Error::invalid_argument("invalid composite key width"));
}
let key = CompositeKey(values.to_vec());
if values.len() == self.column_ids.len() {
let value_to_rows = self.value_to_rows.read();
if let Some(row_ids) = value_to_rows.get(&key) {
buffer.extend_from_slice(row_ids.as_slice());
}
return Ok(());
}
self.ensure_prefix_built(values.len());
let prefix_index = self.prefix_indexes[values.len() - 1].read();
if let Some(row_ids) = prefix_index.get(&key) {
buffer.extend_from_slice(row_ids.as_slice());
}
Ok(())
}
fn get_filtered_row_ids(&self, _expr: &dyn Expression) -> Result<RowIdVec> {
Err(Error::NotSupported(
"multi-column index cannot derive candidates from an arbitrary expression".to_string(),
))
}
fn clear(&self) {
let _publication = self.publication_fence.write();
*self.sorted_values.write() = BTreeMap::new();
self.btree_built.store(false, AtomicOrdering::Release);
*self.value_to_rows.write() = FxHashMap::default();
for (prefix_index, built_flag) in self.prefix_indexes.iter().zip(self.prefix_built.iter()) {
*prefix_index.write() = FxHashMap::default();
built_flag.store(false, AtomicOrdering::Release);
}
*self.row_to_key.write() = I64Map::new();
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
fn close(&mut self) -> Result<()> {
self.closed.store(true, AtomicOrdering::Release);
self.clear();
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::{TimeZone, Utc};
#[test]
fn test_multi_column_index_basic() {
let index = MultiColumnIndex::new(
"test_idx".to_string(),
"test_table".to_string(),
vec!["col1".to_string(), "col2".to_string()],
vec![0, 1],
vec![DataType::Integer, DataType::Text],
false,
0,
);
index
.add(&[Value::Integer(1), Value::Text("a".into())], 100, 0)
.unwrap();
index
.add(&[Value::Integer(1), Value::Text("b".into())], 101, 0)
.unwrap();
index
.add(&[Value::Integer(2), Value::Text("a".into())], 102, 0)
.unwrap();
let results = index
.find(&[Value::Integer(1), Value::Text("a".into())])
.unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0].row_id, 100);
let results = index.find(&[Value::Integer(1)]).unwrap();
assert_eq!(results.len(), 2);
}
#[test]
fn test_multi_column_index_range() {
let index = MultiColumnIndex::new(
"test_idx".to_string(),
"test_table".to_string(),
vec!["amount".to_string()],
vec![0],
vec![DataType::Integer],
false,
0,
);
for i in 0..100 {
index.add(&[Value::Integer(i)], i, 0).unwrap();
}
let results = index
.find_range(&[Value::Integer(10)], &[Value::Integer(20)], true, true)
.unwrap();
assert_eq!(results.len(), 11); }
#[test]
fn leading_timestamp_range_includes_complete_composite_prefix() {
let index = MultiColumnIndex::new(
"due_idx".to_string(),
"jobs".to_string(),
vec!["due_at".to_string(), "id".to_string()],
vec![0, 1],
vec![DataType::Timestamp, DataType::Integer],
false,
0,
);
let first = Utc.with_ymd_and_hms(2026, 8, 9, 0, 0, 0).unwrap();
let second = Utc.with_ymd_and_hms(2026, 8, 10, 0, 0, 0).unwrap();
index
.add(&[Value::Timestamp(first), Value::Integer(1)], 1, 0)
.unwrap();
index
.add(&[Value::Timestamp(first), Value::Integer(2)], 2, 0)
.unwrap();
index
.add(&[Value::Timestamp(second), Value::Integer(3)], 3, 0)
.unwrap();
let inclusive = index
.find_range(&[], &[Value::Timestamp(first)], false, true)
.unwrap();
assert_eq!(
inclusive
.iter()
.map(|entry| entry.row_id)
.collect::<Vec<_>>(),
vec![1, 2]
);
assert!(index
.find_range(&[], &[Value::Timestamp(first)], false, false)
.unwrap()
.is_empty());
assert_eq!(
index
.find_range(&[Value::Timestamp(first)], &[], false, false)
.unwrap()
.iter()
.map(|entry| entry.row_id)
.collect::<Vec<_>>(),
vec![3]
);
}
#[test]
fn partial_composite_bounds_use_lexicographic_ordering() {
let index = MultiColumnIndex::new(
"pair_idx".to_string(),
"pairs".to_string(),
vec!["a".to_string(), "b".to_string(), "id".to_string()],
vec![0, 1, 2],
vec![DataType::Integer; 3],
false,
0,
);
for (row_id, a, b) in [(1, 4, 100), (2, 5, 9), (3, 5, 10), (4, 5, 11)] {
index
.add(
&[Value::Integer(a), Value::Integer(b), Value::Integer(row_id)],
row_id,
0,
)
.unwrap();
}
let rows = index
.find_range(&[], &[Value::Integer(5), Value::Integer(10)], false, true)
.unwrap()
.into_iter()
.map(|entry| entry.row_id)
.collect::<Vec<_>>();
assert_eq!(rows, vec![1, 2, 3]);
}
#[test]
fn test_multi_column_index_unique() {
let index = MultiColumnIndex::new(
"test_idx".to_string(),
"test_table".to_string(),
vec!["id".to_string()],
vec![0],
vec![DataType::Integer],
true, 0,
);
index.add(&[Value::Integer(1)], 100, 0).unwrap();
let result = index.add(&[Value::Integer(1)], 101, 0);
assert!(result.is_err());
}
#[test]
fn test_multi_column_index_remove() {
let index = MultiColumnIndex::new(
"test_idx".to_string(),
"test_table".to_string(),
vec!["col1".to_string()],
vec![0],
vec![DataType::Integer],
false,
0,
);
index.add(&[Value::Integer(1)], 100, 0).unwrap();
let results = index.find(&[Value::Integer(1)]).unwrap();
assert_eq!(results.len(), 1);
index.remove(&[Value::Integer(1)], 100, 0).unwrap();
let results = index.find(&[Value::Integer(1)]).unwrap();
assert_eq!(results.len(), 0);
}
#[test]
fn test_composite_key_ordering() {
let k1 = CompositeKey(vec![Value::Integer(1), Value::Integer(2)]);
let k2 = CompositeKey(vec![Value::Integer(1), Value::Integer(3)]);
let k3 = CompositeKey(vec![Value::Integer(2), Value::Integer(1)]);
assert!(k1 < k2);
assert!(k2 < k3);
assert!(k1 < k3);
}
#[test]
fn v2_r3_multi_authorities_publish_and_reject_stale_remove_together() {
let index = MultiColumnIndex::new(
"idx".into(),
"items".into(),
vec!["a".into(), "b".into()],
vec![0, 1],
vec![DataType::Integer, DataType::Text],
false,
0,
);
let old = [Value::Integer(1), Value::text("old")];
let other = [Value::Integer(2), Value::text("other")];
let stale = [Value::Integer(1), Value::text("stale")];
let replacement = [Value::Integer(3), Value::text("new")];
index.add(&old, 1, 1).unwrap();
index.add(&other, 2, 2).unwrap();
assert_eq!(index.find(&old[..1]).unwrap()[0].row_id, 1);
assert_eq!(index.find_range(&old, &other, true, true).unwrap().len(), 2);
assert!(index.remove(&stale, 1, 1).is_err());
let removals = [(1, old.as_slice()), (2, stale.as_slice())];
assert!(index.remove_batch_slice(&removals).is_err());
assert_eq!(index.find(&old).unwrap()[0].row_id, 1);
assert_eq!(index.find(&old[..1]).unwrap()[0].row_id, 1);
assert_eq!(index.find_range(&old, &other, true, true).unwrap().len(), 2);
index.add(&replacement, 1, 1).unwrap();
assert!(index.find(&old).unwrap().is_empty());
assert!(index.find(&old[..1]).unwrap().is_empty());
assert_eq!(index.find(&replacement).unwrap()[0].row_id, 1);
assert_eq!(index.find(&replacement[..1]).unwrap()[0].row_id, 1);
}
#[test]
fn v2_r3_multi_contradictory_ranges_are_empty() {
let index = MultiColumnIndex::new(
"idx".into(),
"items".into(),
vec!["a".into(), "b".into()],
vec![0, 1],
vec![DataType::Integer, DataType::Integer],
false,
0,
);
index
.add(&[Value::Integer(5), Value::Integer(5)], 1, 1)
.unwrap();
let high = [Value::Integer(10), Value::Integer(0)];
let low = [Value::Integer(1), Value::Integer(0)];
assert!(index
.find_range(&high, &low, true, true)
.unwrap()
.is_empty());
assert!(index
.find_range(&high, &high, false, true)
.unwrap()
.is_empty());
assert!(index
.find_range_ordered_limited(&high, &low, true, true, true, 1)
.unwrap()
.is_empty());
}
}