use std::collections::BTreeMap;
use std::num::NonZeroUsize;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, RwLock};
use ahash::AHashMap;
use lru::LruCache;
use parking_lot::Mutex;
use roaring::RoaringTreemap;
use crate::analysis::analyzer::analyzer::Analyzer;
use crate::analysis::analyzer::standard::StandardAnalyzer;
use crate::analysis::token::Token;
use crate::error::{LaurusError, Result};
use crate::lexical::core::document::Document;
use crate::lexical::core::field::FieldValue;
use crate::lexical::index::inverted::core::posting::{DecodedPostingList, Posting, PostingList};
use crate::lexical::index::inverted::core::terms::{
InvertedIndexTerms, MergedInvertedIndexTerms, TermDictionaryAccess, Terms,
};
use crate::lexical::index::inverted::posting_cache::PostingCache;
use crate::lexical::index::inverted::query_cache::QueryFilterCache;
use crate::lexical::index::inverted::segment::SegmentInfo;
use crate::lexical::index::structures::bkd_tree::{BKDReader, BKDTree};
use crate::lexical::index::structures::dictionary::BlockTermDictionary;
use crate::lexical::index::structures::dictionary::TermInfo;
use crate::lexical::index::structures::doc_values::DocValuesReader;
use crate::lexical::query::Query;
use crate::lexical::reader::FieldStats;
use crate::lexical::reader::PostingIterator;
use crate::maintenance::deletion::DeletionBitmap;
use crate::storage::Storage;
use crate::storage::structured::StructReader;
#[derive(Clone)]
pub struct InvertedIndexReaderConfig {
pub max_cache_memory: usize,
pub enable_term_cache: bool,
pub enable_posting_cache: bool,
pub preload_segments: bool,
pub max_cached_terms_per_field: usize,
pub query_filter_cache_capacity: usize,
pub analyzer: Arc<dyn Analyzer>,
}
impl std::fmt::Debug for InvertedIndexReaderConfig {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("InvertedIndexReaderConfig")
.field("max_cache_memory", &self.max_cache_memory)
.field("enable_term_cache", &self.enable_term_cache)
.field("enable_posting_cache", &self.enable_posting_cache)
.field("preload_segments", &self.preload_segments)
.field(
"max_cached_terms_per_field",
&self.max_cached_terms_per_field,
)
.field(
"query_filter_cache_capacity",
&self.query_filter_cache_capacity,
)
.field("analyzer", &self.analyzer.name())
.finish()
}
}
impl Default for InvertedIndexReaderConfig {
fn default() -> Self {
InvertedIndexReaderConfig {
max_cache_memory: 128 * 1024 * 1024, enable_term_cache: true,
enable_posting_cache: true,
preload_segments: false,
max_cached_terms_per_field: 10000,
query_filter_cache_capacity: 1024,
analyzer: Arc::new(
StandardAnalyzer::new().expect("StandardAnalyzer should be creatable"),
),
}
}
}
#[derive(Debug)]
pub struct InvertedIndexPostingIterator {
data: Arc<DecodedPostingList>,
position: usize,
started: bool,
}
type SoaArrays = (Vec<u32>, Vec<u32>, Option<Vec<Option<Vec<u32>>>>);
impl InvertedIndexPostingIterator {
pub fn new(postings: Vec<Posting>) -> Self {
let (doc_ids, frequencies, positions) = Self::soa_from_aos(&postings);
let skip_levels =
crate::lexical::index::inverted::core::posting::build_skip_levels(&doc_ids);
let doc_frequency = doc_ids.len() as u64;
Self::from_decoded_soa(DecodedPostingList {
term: String::new(),
doc_ids,
frequencies,
weights: Vec::new(),
positions,
skip_levels,
total_frequency: 0,
doc_frequency,
})
}
pub fn with_blocks(postings: Vec<Posting>, _block_size: usize) -> Self {
Self::new(postings)
}
pub fn from_decoded_soa(decoded: DecodedPostingList) -> Self {
Self::from_decoded_soa_arc(Arc::new(decoded))
}
pub fn from_decoded_soa_arc(data: Arc<DecodedPostingList>) -> Self {
InvertedIndexPostingIterator {
data,
position: 0,
started: false,
}
}
pub fn from_decoded_soa_with_blocks(decoded: DecodedPostingList, _block_size: usize) -> Self {
Self::from_decoded_soa(decoded)
}
fn soa_from_aos(postings: &[Posting]) -> SoaArrays {
let n = postings.len();
let mut doc_ids = Vec::with_capacity(n);
let mut frequencies = Vec::with_capacity(n);
let any_positions = postings.iter().any(|p| p.positions.is_some());
let mut positions: Option<Vec<Option<Vec<u32>>>> = if any_positions {
Some(Vec::with_capacity(n))
} else {
None
};
for p in postings {
doc_ids.push(p.doc_id as u32);
frequencies.push(p.frequency);
if let Some(out) = positions.as_mut() {
out.push(p.positions.clone());
}
}
(doc_ids, frequencies, positions)
}
fn skip_via_levels(&self, target_u32: u32) -> usize {
use crate::lexical::index::inverted::core::posting::SKIP_INTERVAL;
let n = self.data.doc_ids.len();
let cursor = self.position;
if cursor >= n {
return n;
}
if self.data.skip_levels.is_empty() {
return cursor;
}
let top = self.data.skip_levels.len() - 1;
let mut step = SKIP_INTERVAL.saturating_pow((top + 1) as u32);
let top_lvl = &self.data.skip_levels[top];
let bucket_lo = cursor / step;
if bucket_lo >= top_lvl.len() {
return cursor;
}
let slice = &top_lvl[bucket_lo..];
let local = slice.partition_point(|&x| x < target_u32);
let mut bucket_index = bucket_lo + local;
let mut lower = bucket_index * step;
for level in (0..top).rev() {
step /= SKIP_INTERVAL;
let lvl = &self.data.skip_levels[level];
let parent_lo = bucket_index * SKIP_INTERVAL;
let parent_hi = (parent_lo + SKIP_INTERVAL).min(lvl.len());
let lo = (cursor / step).max(parent_lo);
if lo >= parent_hi {
lower = lower.max(cursor);
bucket_index = parent_hi.saturating_sub(1);
continue;
}
let slice = &lvl[lo..parent_hi];
let local = slice.partition_point(|&x| x < target_u32);
bucket_index = lo + local;
lower = bucket_index * step;
}
lower.max(cursor).min(n)
}
}
impl crate::lexical::reader::PostingIterator for InvertedIndexPostingIterator {
fn doc_id(&self) -> u64 {
if self.position < self.data.doc_ids.len() {
self.data.doc_ids[self.position] as u64
} else {
u64::MAX }
}
fn term_freq(&self) -> u64 {
if self.position < self.data.frequencies.len() {
self.data.frequencies[self.position] as u64
} else {
0
}
}
fn positions(&self) -> Result<Vec<u64>> {
if self.position >= self.data.doc_ids.len() {
return Ok(Vec::new());
}
match &self.data.positions {
Some(per_doc) => match &per_doc[self.position] {
Some(p) => Ok(p.iter().map(|&v| v as u64).collect()),
None => Ok(Vec::new()),
},
None => Ok(Vec::new()),
}
}
fn next(&mut self) -> Result<bool> {
if self.data.doc_ids.is_empty() {
return Ok(false);
}
if !self.started {
self.started = true;
Ok(true)
} else {
self.position += 1;
Ok(self.position < self.data.doc_ids.len())
}
}
fn skip_to(&mut self, target_doc_id: u64) -> Result<bool> {
self.started = true;
let n = self.data.doc_ids.len();
if n == 0 {
return Ok(false);
}
let target_u32 = match u32::try_from(target_doc_id) {
Ok(t) => t,
Err(_) => {
self.position = n;
return Ok(false);
}
};
self.position = self.skip_via_levels(target_u32);
while self.position < n {
if self.data.doc_ids[self.position] >= target_u32 {
return Ok(true);
}
self.position += 1;
}
Ok(false)
}
fn cost(&self) -> u64 {
self.data.doc_ids.len() as u64
}
}
#[derive(Debug)]
pub struct SegmentReader {
info: SegmentInfo,
storage: Arc<dyn Storage>,
term_dictionary: RwLock<Option<Arc<BlockTermDictionary>>>,
stored_documents: RwLock<Option<BTreeMap<u64, Document>>>,
field_lengths: RwLock<Option<BTreeMap<u64, AHashMap<String, u32>>>>,
field_stats: RwLock<Option<AHashMap<String, crate::lexical::reader::FieldStats>>>,
doc_values: RwLock<Option<Arc<DocValuesReader>>>,
deletion_bitmap: RwLock<Option<Arc<DeletionBitmap>>>,
bkd_trees: RwLock<AHashMap<String, Arc<dyn BKDTree>>>,
posting_cache: PostingCache,
loaded: AtomicBool,
}
impl SegmentReader {
pub fn segment_info(&self) -> &SegmentInfo {
&self.info
}
pub fn open(info: SegmentInfo, storage: Arc<dyn Storage>) -> Result<Self> {
let reader = SegmentReader {
info,
storage,
term_dictionary: RwLock::new(None),
stored_documents: RwLock::new(None),
field_lengths: RwLock::new(None),
field_stats: RwLock::new(None),
doc_values: RwLock::new(None),
deletion_bitmap: RwLock::new(None),
bkd_trees: RwLock::new(AHashMap::new()),
posting_cache: PostingCache::new(0),
loaded: AtomicBool::new(false),
};
Ok(reader)
}
pub fn with_posting_cache_bytes(mut self, max_bytes: usize) -> Self {
self.posting_cache = PostingCache::new(max_bytes);
self
}
pub fn posting_cache_stats(
&self,
) -> crate::lexical::index::inverted::posting_cache::PostingCacheStats {
self.posting_cache.stats()
}
pub fn doc_ids(&self) -> Result<Vec<u64>> {
if !self.loaded.load(Ordering::Acquire) {
self.load_stored_documents()?;
}
let docs = self.stored_documents.read().unwrap();
if let Some(ref documents) = *docs {
Ok(documents.keys().cloned().collect())
} else {
Ok(Vec::new())
}
}
#[deprecated(
since = "0.2.0",
note = "Use `open()` instead. Schema is no longer required."
)]
pub fn open_with_schema(
info: SegmentInfo,
_schema: Arc<()>,
storage: Arc<dyn Storage>,
) -> Result<Self> {
Self::open(info, storage)
}
pub fn load(&mut self) -> Result<()> {
if self.loaded.load(Ordering::Acquire) {
return Ok(());
}
self.load_term_dictionary()?;
self.load_stored_documents()?;
self.load_doc_values()?;
self.load_deletion_bitmap()?;
self.loaded.store(true, Ordering::Release);
Ok(())
}
fn load_term_dictionary(&self) -> Result<()> {
let dict_file = format!("{}.dict", self.info.segment_id);
if let Ok(input) = self.storage.open_input(&dict_file) {
let mut reader = StructReader::new(input)?;
let dictionary = BlockTermDictionary::read_from_storage(&mut reader).map_err(|e| {
LaurusError::index(format!(
"Failed to read term dictionary from {dict_file}: {e}"
))
})?;
*self.term_dictionary.write().unwrap() = Some(Arc::new(dictionary));
}
Ok(())
}
fn load_stored_documents(&self) -> Result<()> {
let docs_file = format!("{}.docs", self.info.segment_id);
if let Ok(input) = self.storage.open_input(&docs_file) {
let mut reader = StructReader::new(input)?;
let doc_count = reader.read_varint()? as usize;
let mut documents = BTreeMap::new();
for _ in 0..doc_count {
let doc_id = reader.read_u64()?;
let field_count = reader.read_varint()? as usize;
let mut doc = Document::new();
for _ in 0..field_count {
let field_name = reader.read_string()?;
let type_tag = reader.read_u8()?;
let field_value = match type_tag {
0 => {
let text = reader.read_string()?;
FieldValue::Text(text)
}
1 => {
let num = reader.read_u64()? as i64;
FieldValue::Int64(num)
}
2 => {
let num = reader.read_f64()?;
FieldValue::Float64(num)
}
3 => {
let b = reader.read_u8()? != 0;
FieldValue::Bool(b)
}
4 => {
let mime = reader.read_string()?;
let data = reader.read_bytes()?;
FieldValue::Bytes(data, if mime.is_empty() { None } else { Some(mime) })
}
5 => {
let dt_str = reader.read_string()?;
let dt = chrono::DateTime::parse_from_rfc3339(&dt_str)
.map_err(|e| {
LaurusError::index(format!("Failed to parse DateTime: {e}"))
})?
.with_timezone(&chrono::Utc);
FieldValue::DateTime(dt)
}
6 => {
let lat = reader.read_f64()?;
let lon = reader.read_f64()?;
FieldValue::Geo(crate::data::GeoPoint::new(lat, lon))
}
7 => {
FieldValue::Null
}
10 => {
let len = reader.read_varint()? as usize;
let mut arr = Vec::with_capacity(len);
for _ in 0..len {
arr.push(reader.read_u64()? as i64);
}
FieldValue::Int64Array(arr)
}
11 => {
let len = reader.read_varint()? as usize;
let mut arr = Vec::with_capacity(len);
for _ in 0..len {
arr.push(reader.read_f64()?);
}
FieldValue::Float64Array(arr)
}
12 => {
let x = reader.read_f64()?;
let y = reader.read_f64()?;
let z = reader.read_f64()?;
FieldValue::GeoEcef(crate::data::GeoEcefPoint::new(x, y, z))
}
_ => {
return Err(LaurusError::index(format!(
"Unknown field type tag: {type_tag}"
)));
}
};
doc.fields.insert(field_name, field_value);
}
documents.insert(doc_id, doc);
}
*self.stored_documents.write().unwrap() = Some(documents);
return Ok(());
}
let json_file = format!("{}.json", self.info.segment_id);
if self.storage.file_exists(&json_file) {
let mut input = self.storage.open_input(&json_file)?;
let mut json_data = String::new();
std::io::Read::read_to_string(&mut input, &mut json_data)?;
let docs: Vec<Document> = serde_json::from_str(&json_data)
.map_err(|e| LaurusError::index(format!("Failed to parse JSON documents: {e}")))?;
let mut documents = BTreeMap::new();
for (idx, doc) in docs.into_iter().enumerate() {
let doc_id = self.info.min_doc_id + idx as u64;
documents.insert(doc_id, doc);
}
*self.stored_documents.write().unwrap() = Some(documents);
}
Ok(())
}
fn load_doc_values(&self) -> Result<()> {
let reader = DocValuesReader::load(self.storage.clone(), &self.info.segment_id)?;
let mut doc_values = self.doc_values.write().unwrap();
*doc_values = Some(Arc::new(reader));
Ok(())
}
fn load_deletion_bitmap(&self) -> Result<()> {
if !self.info.has_deletions {
return Ok(());
}
if self.deletion_bitmap.read().unwrap().is_some() {
return Ok(());
}
let bitmap_file = format!("{}.delmap", self.info.segment_id);
if !self.storage.file_exists(&bitmap_file) {
return Ok(());
}
let input = self.storage.open_input(&bitmap_file)?;
let mut reader = StructReader::new(input)?;
let bitmap = DeletionBitmap::read_from_storage(&mut reader)?;
*self.deletion_bitmap.write().unwrap() = Some(Arc::new(bitmap));
Ok(())
}
pub fn is_deleted(&self, doc_id: u64) -> Result<bool> {
if !self.info.has_deletions {
return Ok(false);
}
if self.deletion_bitmap.read().unwrap().is_none() {
self.load_deletion_bitmap()?;
}
let bitmap_lock = self.deletion_bitmap.read().unwrap();
if let Some(ref bitmap) = *bitmap_lock {
Ok(bitmap.is_deleted(doc_id))
} else {
Ok(false)
}
}
fn filter_deleted_soa(&self, decoded: DecodedPostingList) -> Result<DecodedPostingList> {
if !self.info.has_deletions {
return Ok(decoded);
}
if self.deletion_bitmap.read().unwrap().is_none() {
self.load_deletion_bitmap()?;
}
let bitmap_lock = self.deletion_bitmap.read().unwrap();
let bitmap = match bitmap_lock.as_ref() {
Some(b) => b,
None => return Ok(decoded),
};
let n = decoded.doc_ids.len();
let mut doc_ids = Vec::with_capacity(n);
let mut frequencies = Vec::with_capacity(n);
let mut positions: Option<Vec<Option<Vec<u32>>>> =
decoded.positions.as_ref().map(|_| Vec::with_capacity(n));
for i in 0..n {
let did = decoded.doc_ids[i] as u64;
if bitmap.is_deleted(did) {
continue;
}
doc_ids.push(decoded.doc_ids[i]);
frequencies.push(decoded.frequencies[i]);
if let (Some(out), Some(src)) = (positions.as_mut(), decoded.positions.as_ref()) {
out.push(src[i].clone());
}
}
let skip_levels =
crate::lexical::index::inverted::core::posting::build_skip_levels(&doc_ids);
Ok(DecodedPostingList {
term: decoded.term,
doc_ids,
frequencies,
weights: Vec::new(), positions,
skip_levels,
total_frequency: decoded.total_frequency,
doc_frequency: decoded.doc_frequency,
})
}
fn get_doc_value(&self, field: &str, doc_id: u64) -> Result<Option<FieldValue>> {
if !self.loaded.load(Ordering::Acquire) {
self.ensure_loaded()?;
}
let doc_values = self.doc_values.read().unwrap();
if let Some(reader) = doc_values.as_ref() {
Ok(reader.get_value(field, doc_id).cloned())
} else {
Ok(None)
}
}
fn ensure_loaded(&self) -> Result<()> {
if !self.loaded.load(Ordering::Acquire) {
self.load_doc_values()?;
}
Ok(())
}
fn has_doc_values(&self, field: &str) -> bool {
let doc_values = self.doc_values.read().unwrap();
if let Some(reader) = doc_values.as_ref() {
reader.has_field(field)
} else {
false
}
}
fn load_field_lengths(&self) -> Result<()> {
let lens_file = format!("{}.lens", self.info.segment_id);
if !self.storage.file_exists(&lens_file) {
*self.field_lengths.write().unwrap() = Some(BTreeMap::new());
return Ok(());
}
let lens_input = self.storage.open_input(&lens_file)?;
let mut lens_reader = StructReader::new(lens_input)?;
let doc_count = lens_reader.read_varint()? as usize;
let mut all_field_lengths = BTreeMap::new();
for _ in 0..doc_count {
let doc_id = lens_reader.read_u64()?;
let field_count = lens_reader.read_varint()? as usize;
let mut field_lens = AHashMap::new();
for _ in 0..field_count {
let field_name = lens_reader.read_string()?;
let length = lens_reader.read_u32()?;
field_lens.insert(field_name, length);
}
all_field_lengths.insert(doc_id, field_lens);
}
*self.field_lengths.write().unwrap() = Some(all_field_lengths);
Ok(())
}
fn load_field_stats(&self) -> Result<()> {
let fstats_file = format!("{}.fstats", self.info.segment_id);
if !self.storage.file_exists(&fstats_file) {
*self.field_stats.write().unwrap() = Some(AHashMap::new());
return Ok(());
}
let fstats_input = self.storage.open_input(&fstats_file)?;
let mut fstats_reader = StructReader::new(fstats_input)?;
let field_count = fstats_reader.read_varint()? as usize;
let mut all_field_stats = AHashMap::new();
for _ in 0..field_count {
let field_name = fstats_reader.read_string()?;
let doc_count = fstats_reader.read_u64()?;
let avg_length = fstats_reader.read_f64()?;
let min_length = fstats_reader.read_u64()?;
let max_length = fstats_reader.read_u64()?;
all_field_stats.insert(
field_name.clone(),
crate::lexical::reader::FieldStats {
field: field_name,
unique_terms: 0, total_terms: 0, doc_count,
avg_length,
min_length,
max_length,
},
);
}
*self.field_stats.write().unwrap() = Some(all_field_stats);
Ok(())
}
pub fn field_stats(&self, field: &str) -> Result<Option<FieldStats>> {
if self.field_stats.read().unwrap().is_none() {
self.load_field_stats()?;
}
let field_stats = self.field_stats.read().unwrap();
if let Some(ref stats_map) = *field_stats {
return Ok(stats_map.get(field).cloned());
}
Ok(None)
}
pub fn field_length(&self, doc_id: u64, field: &str) -> Result<Option<u32>> {
if self.is_deleted(doc_id)? {
return Ok(None);
}
let field_lengths = self.field_lengths.read().unwrap();
if let Some(ref lengths_map) = *field_lengths {
return Ok(lengths_map
.get(&doc_id)
.and_then(|doc_lengths| doc_lengths.get(field).copied()));
}
drop(field_lengths);
self.load_field_lengths()?;
let field_lengths = self.field_lengths.read().unwrap();
if let Some(ref lengths_map) = *field_lengths {
return Ok(lengths_map
.get(&doc_id)
.and_then(|doc_lengths| doc_lengths.get(field).copied()));
}
Ok(None)
}
pub fn document(&self, doc_id: u64) -> Result<Option<Document>> {
if !self.loaded.load(Ordering::Acquire) {
self.load_stored_documents()?;
}
if self.is_deleted(doc_id)? {
return Ok(None);
}
let docs = self.stored_documents.read().unwrap();
if let Some(ref documents) = *docs {
Ok(documents.get(&doc_id).cloned())
} else {
Ok(None)
}
}
pub fn document_fields(
&self,
doc_id: u64,
field_names: &[&str],
) -> Result<Option<std::collections::HashMap<String, crate::data::DataValue>>> {
if !self.loaded.load(Ordering::Acquire) {
self.load_stored_documents()?;
}
if self.is_deleted(doc_id)? {
return Ok(None);
}
let docs = self.stored_documents.read().unwrap();
if let Some(ref documents) = *docs
&& let Some(doc) = documents.get(&doc_id)
{
let mut out = std::collections::HashMap::with_capacity(field_names.len());
for &name in field_names {
if let Some(value) = doc.fields.get(name) {
out.insert(name.to_string(), value.clone());
}
}
return Ok(Some(out));
}
Ok(None)
}
pub fn term_info(&self, field: &str, term: &str) -> Result<Option<TermInfo>> {
if self.term_dictionary.read().unwrap().is_none() && !self.loaded.load(Ordering::Acquire) {
self.load_term_dictionary()?;
}
if let Some(ref dict) = *self.term_dictionary.read().unwrap() {
let full_term = format!("{field}:{term}");
Ok(dict.get(&full_term).cloned())
} else {
Ok(None)
}
}
pub fn term_dictionary(&self) -> Result<Option<Arc<BlockTermDictionary>>> {
self.load_term_dictionary()?;
Ok(self.term_dictionary.read().unwrap().clone())
}
pub fn postings(&self, field: &str, term: &str) -> Result<Option<Box<dyn PostingIterator>>> {
let postings_file = format!("{}.post", self.info.segment_id);
if !self.storage.file_exists(&postings_file) {
return self.scan_documents_for_term(field, term);
}
let cache_key = self
.posting_cache
.is_enabled()
.then(|| format!("{field}\u{1}{term}"));
if let Some(key) = &cache_key
&& let Some(cached) = self.posting_cache.get(key)
{
return Ok(Some(Box::new(
InvertedIndexPostingIterator::from_decoded_soa_arc(cached),
)));
}
if let Some(term_info) = self.term_info(field, term)? {
let input = self.storage.open_input(&postings_file)?;
let mut reader = StructReader::new(input)?;
if term_info.posting_offset > 0 {
reader.seek(std::io::SeekFrom::Start(term_info.posting_offset))?;
}
let posting_format = self
.term_dictionary
.read()
.unwrap()
.as_ref()
.map(|dict| dict.posting_format_version())
.unwrap_or(2);
let decoded = if posting_format >= 2 {
PostingList::decode_soa_v2(&mut reader)?
} else {
PostingList::decode_soa(&mut reader)?
};
let filtered = self.filter_deleted_soa(decoded)?;
if filtered.is_empty() {
Ok(None)
} else if let Some(key) = cache_key {
let shared = Arc::new(filtered);
self.posting_cache.put(key, Arc::clone(&shared));
Ok(Some(Box::new(
InvertedIndexPostingIterator::from_decoded_soa_arc(shared),
)))
} else {
Ok(Some(Box::new(
InvertedIndexPostingIterator::from_decoded_soa(filtered),
)))
}
} else {
Ok(None)
}
}
fn scan_documents_for_term(
&self,
field: &str,
term: &str,
) -> Result<Option<Box<dyn PostingIterator>>> {
if !self.loaded.load(Ordering::Acquire) {
self.load_stored_documents()?;
}
let docs = self.stored_documents.read().unwrap();
if let Some(ref documents) = *docs {
let mut postings = Vec::new();
let default_analyzer = StandardAnalyzer::new()?;
for (doc_id, doc) in documents.iter() {
if self.is_deleted(*doc_id)? {
continue;
}
if let Some(field_value) = doc.get_field(field)
&& let Some(text) = field_value.as_text()
{
let token_stream = default_analyzer.analyze(text)?;
let tokens: Vec<Token> = token_stream.collect();
let mut positions = Vec::new();
for token in tokens.iter() {
if token.text == term {
positions.push(token.position as u32);
}
}
if !positions.is_empty() {
postings.push(Posting {
doc_id: *doc_id,
frequency: positions.len() as u32,
positions: Some(positions),
weight: 1.0,
});
}
}
}
if postings.is_empty() {
Ok(None)
} else {
Ok(Some(Box::new(InvertedIndexPostingIterator::with_blocks(
postings, 64,
))))
}
} else {
Ok(None)
}
}
pub fn doc_count(&self) -> u64 {
if !self.info.has_deletions {
return self.info.doc_count;
}
if let Some(bitmap) = self.deletion_bitmap.read().unwrap().clone() {
return bitmap.live_count();
}
if self.load_deletion_bitmap().is_ok()
&& let Some(bitmap) = self.deletion_bitmap.read().unwrap().clone()
{
return bitmap.live_count();
}
self.info.doc_count
}
pub fn get_bkd_tree(&self, field: &str) -> Result<Option<Arc<dyn BKDTree>>> {
if let Some(tree) = self.bkd_trees.read().unwrap().get(field) {
return Ok(Some(tree.clone()));
}
let bkd_file = format!("{}.{}.bkd", self.info.segment_id, field);
if self.storage.file_exists(&bkd_file) {
let reader = BKDReader::open(self.storage.clone(), &bkd_file)?;
let tree: Arc<dyn BKDTree> = Arc::new(reader);
self.bkd_trees
.write()
.unwrap()
.insert(field.to_string(), tree.clone());
return Ok(Some(tree));
}
Ok(None)
}
pub(crate) fn get_filtered_bkd_tree(&self, field: &str) -> Result<Option<Arc<dyn BKDTree>>> {
let Some(tree) = self.get_bkd_tree(field)? else {
return Ok(None);
};
if !self.info.has_deletions {
return Ok(Some(tree));
}
self.load_deletion_bitmap()?;
let Some(bitmap) = self.deletion_bitmap.read().unwrap().clone() else {
return Ok(Some(tree));
};
let snapshot = Arc::new(DeletionSnapshot {
bitmaps: vec![(self.info.min_doc_id, self.info.max_doc_id, bitmap)],
});
Ok(Some(Arc::new(DeletionFilteringBKDTree {
inner: tree,
snapshot,
})))
}
}
#[derive(Debug)]
struct MultiSegmentBKDTree {
trees: Vec<Arc<dyn BKDTree>>,
}
impl BKDTree for MultiSegmentBKDTree {
fn intersect(
&self,
visitor: &mut dyn crate::lexical::index::structures::visitor::IntersectVisitor,
) -> Result<()> {
for tree in &self.trees {
tree.intersect(visitor)?;
}
Ok(())
}
}
#[derive(Debug, Clone)]
struct DeletionSnapshot {
bitmaps: Vec<(u64, u64, Arc<DeletionBitmap>)>,
}
impl DeletionSnapshot {
#[inline]
fn is_empty(&self) -> bool {
self.bitmaps.is_empty()
}
#[inline]
fn is_deleted(&self, doc_id: u64) -> bool {
for (min, max, bitmap) in &self.bitmaps {
if doc_id >= *min && doc_id <= *max && bitmap.is_deleted(doc_id) {
return true;
}
}
false
}
}
struct DeletionFilteringBKDTree {
inner: Arc<dyn BKDTree>,
snapshot: Arc<DeletionSnapshot>,
}
impl std::fmt::Debug for DeletionFilteringBKDTree {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("DeletionFilteringBKDTree")
.field("inner", &self.inner)
.field(
"snapshot_segments_with_deletions",
&self.snapshot.bitmaps.len(),
)
.finish()
}
}
impl BKDTree for DeletionFilteringBKDTree {
fn intersect(
&self,
visitor: &mut dyn crate::lexical::index::structures::visitor::IntersectVisitor,
) -> Result<()> {
if self.snapshot.is_empty() {
return self.inner.intersect(visitor);
}
let mut wrapped = DeletionFilteringVisitor {
inner: visitor,
snapshot: &self.snapshot,
};
self.inner.intersect(&mut wrapped)
}
}
struct DeletionFilteringVisitor<'a> {
inner: &'a mut dyn crate::lexical::index::structures::visitor::IntersectVisitor,
snapshot: &'a DeletionSnapshot,
}
impl crate::lexical::index::structures::visitor::IntersectVisitor for DeletionFilteringVisitor<'_> {
fn compare(
&self,
cell: &crate::lexical::index::structures::aabb::AABB,
) -> crate::lexical::index::structures::visitor::CellRelation {
self.inner.compare(cell)
}
fn visit_inside(&mut self, doc_id: u64) {
if !self.snapshot.is_deleted(doc_id) {
self.inner.visit_inside(doc_id);
}
}
fn visit(&mut self, doc_id: u64, point: &[f64]) {
if !self.snapshot.is_deleted(doc_id) {
self.inner.visit(doc_id, point);
}
}
}
const EST_TERM_ENTRY_BYTES: usize = 64;
#[derive(Debug)]
pub struct CacheManager {
term_cache: Mutex<LruCache<String, Arc<TermInfo>>>,
memory_limit: usize,
cache_hits: AtomicUsize,
cache_misses: AtomicUsize,
}
impl CacheManager {
pub fn new(memory_limit: usize) -> Self {
let capacity = NonZeroUsize::new((memory_limit / EST_TERM_ENTRY_BYTES).max(1))
.unwrap_or(NonZeroUsize::MIN);
CacheManager {
term_cache: Mutex::new(LruCache::new(capacity)),
memory_limit,
cache_hits: AtomicUsize::new(0),
cache_misses: AtomicUsize::new(0),
}
}
pub fn get_term_info(&self, key: &str) -> Option<Arc<TermInfo>> {
let hit = self.term_cache.lock().get(key).cloned();
if hit.is_some() {
self.cache_hits.fetch_add(1, Ordering::Relaxed);
} else {
self.cache_misses.fetch_add(1, Ordering::Relaxed);
}
hit
}
pub fn cache_term_info(&self, key: String, info: TermInfo) {
self.term_cache.lock().put(key, Arc::new(info));
}
pub fn stats(&self) -> CacheStats {
let entries = self.term_cache.lock().len();
CacheStats {
hits: self.cache_hits.load(Ordering::Relaxed),
misses: self.cache_misses.load(Ordering::Relaxed),
memory_usage: entries * EST_TERM_ENTRY_BYTES,
memory_limit: self.memory_limit,
}
}
}
#[derive(Debug, Clone)]
pub struct CacheStats {
pub hits: usize,
pub misses: usize,
pub memory_usage: usize,
pub memory_limit: usize,
}
impl CacheStats {
pub fn hit_ratio(&self) -> f64 {
if self.hits + self.misses == 0 {
0.0
} else {
self.hits as f64 / (self.hits + self.misses) as f64
}
}
}
#[derive(Debug, Clone)]
pub struct InvertedIndexReader {
segment_readers: Vec<Arc<RwLock<SegmentReader>>>,
segment_infos: Vec<SegmentInfo>,
cache_manager: Arc<CacheManager>,
query_cache: Arc<QueryFilterCache>,
config: InvertedIndexReaderConfig,
closed: Arc<AtomicBool>,
total_doc_count: u64,
}
impl InvertedIndexReader {
pub fn new(
segments: Vec<SegmentInfo>,
storage: Arc<dyn Storage>,
config: InvertedIndexReaderConfig,
) -> Result<Self> {
let cache_manager = Arc::new(CacheManager::new(config.max_cache_memory));
let query_cache = Arc::new(QueryFilterCache::new(config.query_filter_cache_capacity));
let mut segment_readers = Vec::new();
let mut total_doc_count = 0;
let posting_cache_bytes = if config.enable_posting_cache {
config.max_cache_memory
} else {
0
};
for segment_info in &segments {
total_doc_count += segment_info.doc_count;
let mut reader = SegmentReader::open(segment_info.clone(), storage.clone())?
.with_posting_cache_bytes(posting_cache_bytes);
if config.preload_segments {
reader.load()?;
}
segment_readers.push(Arc::new(RwLock::new(reader)));
}
Ok(InvertedIndexReader {
segment_readers,
segment_infos: segments,
cache_manager,
query_cache,
config,
closed: Arc::new(AtomicBool::new(false)),
total_doc_count,
})
}
pub fn cache_stats(&self) -> CacheStats {
self.cache_manager.stats()
}
pub fn query_cache_stats(
&self,
) -> crate::lexical::index::inverted::query_cache::QueryFilterCacheStats {
self.query_cache.stats()
}
pub fn matching_doc_ids(&self, query: &dyn Query) -> Result<Arc<RoaringTreemap>> {
if let Some(key) = query.cache_key() {
if let Some(cached) = self.query_cache.get(&key) {
return Ok(cached);
}
let bitmap = Arc::new(self.drain_matching(query)?);
self.query_cache.put(key, Arc::clone(&bitmap));
Ok(bitmap)
} else {
Ok(Arc::new(self.drain_matching(query)?))
}
}
fn drain_matching(&self, query: &dyn Query) -> Result<RoaringTreemap> {
let matcher = query.matcher(self)?;
crate::lexical::index::inverted::query_cache::drain_matcher(matcher)
}
pub fn analyzer(&self) -> &Arc<dyn Analyzer> {
&self.config.analyzer
}
pub fn segment_count(&self) -> usize {
self.segment_readers.len()
}
pub fn segment_readers(&self) -> &[Arc<RwLock<SegmentReader>>] {
&self.segment_readers
}
fn check_closed(&self) -> Result<()> {
if self.closed.load(Ordering::Acquire) {
Err(LaurusError::index("Reader is closed"))
} else {
Ok(())
}
}
pub fn field_length(&self, doc_id: u64, field: &str) -> Result<Option<u32>> {
self.check_closed()?;
for (i, segment_reader) in self.segment_readers.iter().enumerate() {
if let Some(info) = self.segment_infos.get(i)
&& (doc_id < info.min_doc_id || doc_id > info.max_doc_id)
{
continue;
}
let reader = segment_reader.read().unwrap();
if let Ok(Some(length)) = reader.field_length(doc_id, field) {
return Ok(Some(length));
}
}
Ok(None)
}
}
impl crate::lexical::reader::LexicalIndexReader for InvertedIndexReader {
fn doc_count(&self) -> u64 {
self.segment_readers
.iter()
.map(|sr| sr.read().unwrap().doc_count())
.sum()
}
fn max_doc(&self) -> u64 {
self.total_doc_count
}
fn is_deleted(&self, doc_id: u64) -> bool {
for segment_reader in &self.segment_readers {
let reader = segment_reader.read().unwrap();
if let Ok(true) = reader.is_deleted(doc_id) {
return true;
}
}
false
}
fn document(&self, doc_id: u64) -> Result<Option<Document>> {
self.check_closed()?;
for segment_reader in &self.segment_readers {
let reader = segment_reader.read().unwrap();
if let Ok(Some(doc)) = reader.document(doc_id) {
return Ok(Some(doc));
}
}
Ok(None)
}
fn document_fields(
&self,
doc_id: u64,
field_names: &[&str],
) -> Result<Option<std::collections::HashMap<String, crate::data::DataValue>>> {
self.check_closed()?;
for segment_reader in &self.segment_readers {
let reader = segment_reader.read().unwrap();
if let Ok(Some(fields)) = reader.document_fields(doc_id, field_names) {
return Ok(Some(fields));
}
}
Ok(None)
}
fn doc_ids(&self) -> Result<Vec<u64>> {
self.check_closed()?;
let mut all_ids = Vec::new();
for segment_reader in &self.segment_readers {
let reader = segment_reader.read().unwrap();
all_ids.extend(reader.doc_ids()?);
}
Ok(all_ids)
}
fn term_info(
&self,
field: &str,
term: &str,
) -> Result<Option<crate::lexical::reader::ReaderTermInfo>> {
self.check_closed()?;
let cache_key = format!("{field}:{term}");
if let Some(cached_info) = self.cache_manager.get_term_info(&cache_key) {
return Ok(Some(crate::lexical::reader::ReaderTermInfo {
field: field.to_string(),
term: term.to_string(),
doc_freq: cached_info.doc_frequency,
total_freq: cached_info.total_frequency,
posting_offset: cached_info.posting_offset,
posting_size: cached_info.posting_length,
max_score_factor: cached_info.max_score_factor,
block_max: cached_info.block_max.clone(),
}));
}
let mut total_doc_freq = 0;
let mut total_term_freq = 0;
let mut max_score_factor: f32 = 0.0;
let mut matched_count = 0_usize;
let mut combined_block_max: Vec<crate::lexical::index::structures::dictionary::BlockMax> =
Vec::new();
for segment_reader in &self.segment_readers {
let reader = segment_reader.read().unwrap();
if let Some(term_info) = reader.term_info(field, term)? {
total_doc_freq += term_info.doc_frequency;
total_term_freq += term_info.total_frequency;
max_score_factor = max_score_factor.max(term_info.max_score_factor);
matched_count += 1;
combined_block_max.extend(term_info.block_max.iter().copied());
}
}
let found = matched_count > 0;
let aggregated_block_max = if matched_count == 1 {
combined_block_max
} else {
Vec::new()
};
if found {
let reader_info = crate::lexical::reader::ReaderTermInfo {
field: field.to_string(),
term: term.to_string(),
doc_freq: total_doc_freq,
total_freq: total_term_freq,
posting_offset: 0, posting_size: 0, max_score_factor,
block_max: aggregated_block_max.clone(),
};
let term_info = TermInfo {
posting_offset: 0,
posting_length: 0,
doc_frequency: total_doc_freq,
total_frequency: total_term_freq,
max_score_factor,
block_max: aggregated_block_max,
};
self.cache_manager.cache_term_info(cache_key, term_info);
Ok(Some(reader_info))
} else {
Ok(None)
}
}
fn postings(
&self,
field: &str,
term: &str,
) -> Result<Option<Box<dyn crate::lexical::reader::PostingIterator>>> {
self.check_closed()?;
let mut iterators = Vec::new();
for segment_reader in &self.segment_readers {
let reader = segment_reader.read().unwrap();
if let Some(iter) = reader.postings(field, term)? {
iterators.push(iter);
}
}
if iterators.is_empty() {
Ok(None)
} else if iterators.len() == 1 {
Ok(Some(iterators.into_iter().next().unwrap()))
} else {
let merged = MergedPostingIterator::new(iterators)?;
Ok(Some(Box::new(merged)))
}
}
fn field_stats(&self, field: &str) -> Result<Option<crate::lexical::reader::FieldStats>> {
self.check_closed()?;
let mut total_doc_count = 0u64;
let mut total_length_sum = 0u64; let mut min_length = u64::MAX;
let mut max_length = 0u64;
let mut found = false;
for segment_reader in &self.segment_readers {
let reader = segment_reader.read().unwrap();
if let Some(segment_stats) = reader.field_stats(field)? {
total_doc_count += segment_stats.doc_count;
total_length_sum +=
(segment_stats.avg_length * segment_stats.doc_count as f64) as u64;
min_length = min_length.min(segment_stats.min_length);
max_length = max_length.max(segment_stats.max_length);
found = true;
}
}
if found {
Ok(Some(crate::lexical::reader::FieldStats {
field: field.to_string(),
unique_terms: 0, total_terms: 0, doc_count: total_doc_count,
avg_length: if total_doc_count > 0 {
total_length_sum as f64 / total_doc_count as f64
} else {
0.0
},
min_length: if min_length == u64::MAX {
0
} else {
min_length
},
max_length,
}))
} else {
Ok(None)
}
}
fn close(&mut self) -> Result<()> {
self.closed.store(true, Ordering::Release);
Ok(())
}
fn is_closed(&self) -> bool {
self.closed.load(Ordering::Acquire)
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
fn get_doc_value(&self, field: &str, doc_id: u64) -> Result<Option<FieldValue>> {
for segment_lock in &self.segment_readers {
let segment = segment_lock.read().unwrap();
if let Ok(Some(value)) = segment.get_doc_value(field, doc_id) {
return Ok(Some(value));
}
}
Ok(None)
}
fn has_doc_values(&self, field: &str) -> bool {
self.segment_readers.iter().any(|seg_lock| {
let seg = seg_lock.read().unwrap();
seg.has_doc_values(field)
})
}
fn get_bkd_tree(&self, field: &str) -> Result<Option<Arc<dyn BKDTree>>> {
self.check_closed()?;
let mut trees = Vec::new();
for segment_reader in &self.segment_readers {
let reader = segment_reader.read().unwrap();
if let Some(tree) = reader.get_bkd_tree(field)? {
trees.push(tree);
}
}
if trees.is_empty() {
return Ok(None);
}
let multi: Arc<dyn BKDTree> = Arc::new(MultiSegmentBKDTree { trees });
let mut bitmaps = Vec::new();
for sr in &self.segment_readers {
let reader = sr.read().unwrap();
if !reader.info.has_deletions {
continue;
}
reader.load_deletion_bitmap()?;
if let Some(bitmap) = reader.deletion_bitmap.read().unwrap().clone() {
bitmaps.push((reader.info.min_doc_id, reader.info.max_doc_id, bitmap));
}
}
let snapshot = Arc::new(DeletionSnapshot { bitmaps });
Ok(Some(Arc::new(DeletionFilteringBKDTree {
inner: multi,
snapshot,
})))
}
}
impl TermDictionaryAccess for InvertedIndexReader {
fn terms(&self, field: &str) -> Result<Option<Box<dyn Terms>>> {
let mut dicts = Vec::new();
for seg_lock in &self.segment_readers {
let seg = seg_lock.read().unwrap();
if seg.term_dictionary.read().unwrap().is_none() {
seg.load_term_dictionary()?;
}
if let Some(dict) = seg.term_dictionary.read().unwrap().clone() {
dicts.push(dict);
}
}
if dicts.is_empty() {
return Ok(None);
}
if dicts.len() == 1 {
let terms = InvertedIndexTerms::new(field, dicts.into_iter().next().unwrap());
return Ok(Some(Box::new(terms)));
}
let terms = MergedInvertedIndexTerms::new(field, &dicts);
Ok(Some(Box::new(terms)))
}
}
#[derive(Debug)]
pub struct MergedPostingIterator {
inner: MergeImpl,
current_doc: u64,
started: bool,
}
const LINEAR_THRESHOLD: usize = 8;
#[derive(Debug)]
enum MergeImpl {
Linear {
wrappers: Vec<IteratorWrapper>,
min_idx: usize,
},
Heap(std::collections::BinaryHeap<IteratorWrapper>),
}
#[derive(Debug)]
struct IteratorWrapper {
iter: Box<dyn crate::lexical::reader::PostingIterator>,
current_doc: u64,
}
impl PartialEq for IteratorWrapper {
fn eq(&self, other: &Self) -> bool {
self.current_doc == other.current_doc
}
}
impl Eq for IteratorWrapper {}
impl PartialOrd for IteratorWrapper {
fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
Some(self.cmp(other))
}
}
impl Ord for IteratorWrapper {
fn cmp(&self, other: &Self) -> std::cmp::Ordering {
other.current_doc.cmp(&self.current_doc)
}
}
#[inline]
fn find_min_idx(wrappers: &[IteratorWrapper]) -> usize {
let mut min_idx = 0;
let mut min_doc = wrappers[0].current_doc;
for (i, w) in wrappers.iter().enumerate().skip(1) {
if w.current_doc < min_doc {
min_doc = w.current_doc;
min_idx = i;
}
}
min_idx
}
impl MergedPostingIterator {
pub fn new(iterators: Vec<Box<dyn crate::lexical::reader::PostingIterator>>) -> Result<Self> {
let mut wrappers = Vec::with_capacity(iterators.len());
for mut iter in iterators {
if iter.next()? {
let doc_id = iter.doc_id();
wrappers.push(IteratorWrapper {
iter,
current_doc: doc_id,
});
}
}
let inner = if wrappers.len() <= LINEAR_THRESHOLD {
let min_idx = if wrappers.is_empty() {
0
} else {
find_min_idx(&wrappers)
};
MergeImpl::Linear { wrappers, min_idx }
} else {
let mut heap = std::collections::BinaryHeap::with_capacity(wrappers.len());
for w in wrappers {
heap.push(w);
}
MergeImpl::Heap(heap)
};
let current_doc = match &inner {
MergeImpl::Linear { wrappers, min_idx } => {
if wrappers.is_empty() {
u64::MAX
} else {
wrappers[*min_idx].current_doc
}
}
MergeImpl::Heap(heap) => heap.peek().map_or(u64::MAX, |w| w.current_doc),
};
Ok(MergedPostingIterator {
inner,
current_doc,
started: false,
})
}
fn advance(&mut self) -> Result<bool> {
match &mut self.inner {
MergeImpl::Linear { wrappers, min_idx } => {
if wrappers.is_empty() {
self.current_doc = u64::MAX;
return Ok(false);
}
let idx = *min_idx;
if wrappers[idx].iter.next()? {
wrappers[idx].current_doc = wrappers[idx].iter.doc_id();
} else {
wrappers.swap_remove(idx);
}
if wrappers.is_empty() {
self.current_doc = u64::MAX;
Ok(false)
} else {
*min_idx = find_min_idx(wrappers);
self.current_doc = wrappers[*min_idx].current_doc;
Ok(true)
}
}
MergeImpl::Heap(heap) => {
if let Some(mut wrapper) = heap.pop() {
if wrapper.iter.next()? {
wrapper.current_doc = wrapper.iter.doc_id();
heap.push(wrapper);
}
if let Some(new_top) = heap.peek() {
self.current_doc = new_top.current_doc;
Ok(true)
} else {
self.current_doc = u64::MAX;
Ok(false)
}
} else {
self.current_doc = u64::MAX;
Ok(false)
}
}
}
}
fn current_wrapper(&self) -> Option<&IteratorWrapper> {
match &self.inner {
MergeImpl::Linear { wrappers, min_idx } => {
if wrappers.is_empty() {
None
} else {
Some(&wrappers[*min_idx])
}
}
MergeImpl::Heap(heap) => heap.peek(),
}
}
}
impl crate::lexical::reader::PostingIterator for MergedPostingIterator {
fn doc_id(&self) -> u64 {
self.current_doc
}
fn term_freq(&self) -> u64 {
self.current_wrapper().map_or(0, |w| w.iter.term_freq())
}
fn positions(&self) -> Result<Vec<u64>> {
self.current_wrapper()
.map_or(Ok(Vec::new()), |w| w.iter.positions())
}
fn next(&mut self) -> Result<bool> {
if !self.started {
self.started = true;
let exhausted = match &self.inner {
MergeImpl::Linear { wrappers, .. } => wrappers.is_empty(),
MergeImpl::Heap(heap) => heap.is_empty(),
};
return Ok(!exhausted);
}
self.advance()
}
fn skip_to(&mut self, target: u64) -> Result<bool> {
if !self.started {
self.started = true;
}
while self.doc_id() < target {
if !self.advance()? {
return Ok(false);
}
}
Ok(self.doc_id() != u64::MAX)
}
fn cost(&self) -> u64 {
match &self.inner {
MergeImpl::Linear { wrappers, .. } => wrappers.iter().map(|w| w.iter.cost()).sum(),
MergeImpl::Heap(heap) => heap.iter().map(|w| w.iter.cost()).sum(),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::lexical::reader::PostingIterator;
#[test]
fn test_advanced_posting_iterator() {
let postings = vec![
crate::lexical::index::inverted::core::posting::Posting {
doc_id: 1,
frequency: 1,
positions: Some(vec![0]),
weight: 1.0,
},
crate::lexical::index::inverted::core::posting::Posting {
doc_id: 3,
frequency: 1,
positions: Some(vec![0]),
weight: 1.0,
},
crate::lexical::index::inverted::core::posting::Posting {
doc_id: 5,
frequency: 1,
positions: Some(vec![0]),
weight: 1.0,
},
crate::lexical::index::inverted::core::posting::Posting {
doc_id: 7,
frequency: 1,
positions: Some(vec![0]),
weight: 1.0,
},
crate::lexical::index::inverted::core::posting::Posting {
doc_id: 9,
frequency: 1,
positions: Some(vec![0]),
weight: 1.0,
},
];
let mut iter = InvertedIndexPostingIterator::with_blocks(postings, 2);
assert!(iter.skip_to(5).unwrap());
assert_eq!(iter.doc_id(), 5);
assert!(iter.next().unwrap());
assert_eq!(iter.doc_id(), 7);
assert!(!iter.skip_to(15).unwrap());
assert_eq!(iter.doc_id(), u64::MAX);
}
#[test]
fn from_decoded_soa_arc_shares_backing_without_deep_clone() {
use crate::lexical::index::inverted::core::posting::{
DecodedPostingList, build_skip_levels,
};
let doc_ids = vec![2u32, 4, 6, 8, 10];
let skip_levels = build_skip_levels(&doc_ids);
let shared = Arc::new(DecodedPostingList {
term: "t".to_string(),
doc_ids: doc_ids.clone(),
frequencies: vec![1, 1, 1, 1, 1],
weights: Vec::new(),
positions: None,
skip_levels,
total_frequency: 5,
doc_frequency: 5,
});
assert_eq!(Arc::strong_count(&shared), 1);
let mut it_a = InvertedIndexPostingIterator::from_decoded_soa_arc(Arc::clone(&shared));
let mut it_b = InvertedIndexPostingIterator::from_decoded_soa_arc(Arc::clone(&shared));
assert_eq!(
Arc::strong_count(&shared),
3,
"both iterators must share the cached Arc (no deep clone)"
);
assert!(it_a.skip_to(6).unwrap());
assert_eq!(it_a.doc_id(), 6);
assert!(it_b.next().unwrap());
assert_eq!(it_b.doc_id(), 2);
let mut seen = vec![it_b.doc_id()];
while it_b.next().unwrap() {
seen.push(it_b.doc_id());
}
assert_eq!(seen, vec![2, 4, 6, 8, 10]);
drop(it_a);
drop(it_b);
assert_eq!(Arc::strong_count(&shared), 1);
}
#[test]
fn test_skip_to_matches_linear_scan() {
use crate::lexical::index::inverted::core::posting::{Posting, SKIP_INTERVAL};
for &n in &[
1usize,
SKIP_INTERVAL - 1,
SKIP_INTERVAL,
SKIP_INTERVAL + 1,
SKIP_INTERVAL * SKIP_INTERVAL,
5_000,
] {
let postings: Vec<Posting> = (0..n as u64)
.map(|i| Posting::with_frequency(i * 3 + 7, 1))
.collect();
let doc_ids: Vec<u64> = postings.iter().map(|p| p.doc_id).collect();
let mut targets: Vec<u64> = vec![0, doc_ids[0]];
if n > 1 {
targets.push(doc_ids[n / 2]);
targets.push(doc_ids[n / 2] + 1);
}
targets.push(doc_ids[n - 1]);
targets.push(doc_ids[n - 1] + 1);
targets.push(doc_ids[n - 1] + 1000);
targets.push(u64::from(u32::MAX) + 1);
for &target in &targets {
let mut iter = InvertedIndexPostingIterator::new(postings.clone());
let got = iter.skip_to(target).unwrap();
let want_idx = doc_ids.iter().position(|&d| d >= target);
match want_idx {
Some(idx) => {
assert!(got, "expected hit at target={target} n={n}");
assert_eq!(
iter.doc_id(),
doc_ids[idx],
"wrong doc_id at target={target} n={n}"
);
}
None => {
assert!(!got, "expected miss at target={target} n={n}");
assert_eq!(
iter.doc_id(),
u64::MAX,
"exhausted iter should report u64::MAX (target={target} n={n})"
);
}
}
}
}
}
#[test]
fn test_skip_to_is_monotonic() {
use crate::lexical::index::inverted::core::posting::Posting;
let n: usize = 2_048;
let postings: Vec<Posting> = (0..n as u64)
.map(|i| Posting::with_frequency(i * 2 + 1, 1))
.collect();
let mut iter = InvertedIndexPostingIterator::new(postings);
let mut prev_doc: u64 = 0;
for target in (50..n as u64 * 2).step_by(101) {
assert!(iter.skip_to(target).unwrap(), "target={target}");
let current = iter.doc_id();
assert!(
current >= prev_doc,
"regressed: prev={prev_doc} current={current} target={target}"
);
assert!(
current >= target,
"landed before target: current={current} target={target}"
);
prev_doc = current;
}
}
#[test]
fn test_cache_manager() {
let cache = CacheManager::new(1024);
let key = "field:term".to_string();
let term_info = TermInfo::new(100, 50, 5, 10);
assert!(cache.get_term_info(&key).is_none());
cache.cache_term_info(key.clone(), term_info.clone());
let cached = cache.get_term_info(&key).unwrap();
assert_eq!(cached.doc_frequency, term_info.doc_frequency);
let stats = cache.stats();
assert_eq!(stats.hits, 1);
assert_eq!(stats.misses, 1);
assert!(stats.hit_ratio() > 0.0);
}
#[test]
fn test_cache_manager_lru_eviction() {
let cache = CacheManager::new(2 * EST_TERM_ENTRY_BYTES);
cache.cache_term_info("a".to_string(), TermInfo::new(1, 1, 1, 1));
cache.cache_term_info("b".to_string(), TermInfo::new(2, 2, 2, 2));
assert!(cache.get_term_info("a").is_some());
cache.cache_term_info("c".to_string(), TermInfo::new(3, 3, 3, 3));
assert!(
cache.get_term_info("a").is_some(),
"recently-used 'a' survives"
);
assert!(
cache.get_term_info("c").is_some(),
"just-inserted 'c' survives"
);
assert!(
cache.get_term_info("b").is_none(),
"least-recently-used 'b' must be evicted, not a random entry"
);
}
#[test]
fn posting_cache_hit_and_snapshot_invalidation() {
use crate::Document;
use crate::lexical::store::LexicalStore;
use crate::lexical::store::config::LexicalIndexConfig;
use crate::storage::memory::{MemoryStorage, MemoryStorageConfig};
let storage = Arc::new(MemoryStorage::new(MemoryStorageConfig::default()));
let store = LexicalStore::new(storage, LexicalIndexConfig::default()).unwrap();
for id in 0..5u64 {
store
.upsert_document(
id,
Document::builder().add_text("body", "shared term").build(),
)
.unwrap();
}
store.commit().unwrap();
let drain = |it: Option<Box<dyn crate::lexical::reader::PostingIterator>>| -> Vec<u64> {
let mut ids = Vec::new();
if let Some(mut it) = it {
while it.next().unwrap() {
ids.push(it.doc_id());
}
}
ids.sort_unstable();
ids
};
{
let reader = store.reader_for_tests().unwrap();
let inverted = reader
.as_any()
.downcast_ref::<InvertedIndexReader>()
.unwrap();
let seg = inverted.segment_readers()[0].read().unwrap();
let first = drain(seg.postings("body", "shared").unwrap());
let second = drain(seg.postings("body", "shared").unwrap());
assert_eq!(first, vec![0, 1, 2, 3, 4]);
assert_eq!(first, second, "cached postings must match the decoded list");
let stats = seg.posting_cache_stats();
assert_eq!(stats.misses, 1, "the first decode is a cache miss");
assert!(stats.hits >= 1, "the repeat lookup hits the cache");
}
store.delete_document_by_internal_id(2).unwrap();
store.commit().unwrap();
let reader2 = store.reader_for_tests().unwrap();
let after = drain(reader2.postings("body", "shared").unwrap());
assert_eq!(
after,
vec![0, 1, 3, 4],
"deleted doc 2 must be excluded in the new snapshot"
);
}
#[test]
fn test_segment_info() {
let info = SegmentInfo {
segment_id: "seg_000001".to_string(),
doc_count: 1000,
min_doc_id: 0,
max_doc_id: 999,
generation: 1,
has_deletions: false,
shard_id: 0,
};
assert_eq!(info.segment_id, "seg_000001");
assert_eq!(info.doc_count, 1000);
assert_eq!(info.min_doc_id, 0);
assert_eq!(info.max_doc_id, 999);
assert!(!info.has_deletions);
}
}