1#[cfg(feature = "native")]
11use crate::dsl::Schema;
12#[cfg(feature = "native")]
13use crate::error::Result;
14#[cfg(feature = "sync")]
15use std::collections::HashMap;
16#[cfg(feature = "native")]
17use std::sync::Arc;
18#[cfg(feature = "sync")]
19use std::sync::{OnceLock, Weak};
20
21mod searcher;
22pub use searcher::Searcher;
23
24#[cfg(feature = "native")]
25mod primary_key;
26#[cfg(feature = "native")]
27mod reader;
28#[cfg(feature = "native")]
29mod vector_builder;
30#[cfg(all(feature = "wasm", not(feature = "native")))]
31mod wasm_writer;
32#[cfg(feature = "native")]
33mod writer;
34#[cfg(feature = "native")]
35pub use primary_key::PrimaryKeyIndex;
36#[cfg(feature = "native")]
37pub use reader::IndexReader;
38#[cfg(all(feature = "wasm", not(feature = "native")))]
39pub use wasm_writer::IndexWriter as WasmIndexWriter;
40#[cfg(feature = "native")]
41pub use writer::{IndexWriter, PreparedCommit, WRITER_LOCK_FILENAME};
42
43mod metadata;
44pub use metadata::{
45 FieldVectorMeta, INDEX_META_FILENAME, IndexMetadata, SegmentMetaInfo, VectorIndexState,
46};
47
48#[cfg(feature = "native")]
49mod helpers;
50#[cfg(feature = "native")]
51pub use helpers::{
52 IndexingStats, SchemaConfig, SchemaFieldConfig, create_index_at_path, create_index_from_sdl,
53 index_documents_from_reader, index_json_document, parse_schema,
54};
55
56pub const SLICE_CACHE_FILENAME: &str = "index.slicecache";
58
59#[derive(Debug, Clone)]
61pub struct IndexConfig {
62 pub num_threads: usize,
68 pub num_indexing_threads: usize,
70 pub num_compression_threads: usize,
72 pub term_cache_blocks: usize,
74 pub store_cache_blocks: usize,
76 pub max_indexing_memory_bytes: usize,
78 pub vector_training_max_samples: usize,
82 pub vector_training_memory_bytes: usize,
84 pub merge_policy: Box<dyn crate::merge::MergePolicy>,
86 pub optimization: crate::structures::IndexOptimization,
88 pub reload_interval_ms: u64,
90 pub max_concurrent_merges: usize,
92 #[cfg(feature = "native")]
96 pub background_merge_permits: Arc<tokio::sync::Semaphore>,
97 pub merge_bp_time_budget: Option<std::time::Duration>,
104 pub bp_memory_budget_bytes: usize,
111 #[cfg(feature = "native")]
117 pub background_reorder_permits: Arc<tokio::sync::Semaphore>,
118 #[cfg(feature = "native")]
122 pub background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
123}
124
125#[cfg(feature = "sync")]
129static SEARCH_CPU_POOLS: OnceLock<parking_lot::Mutex<HashMap<usize, Weak<rayon::ThreadPool>>>> =
130 OnceLock::new();
131
132#[cfg(feature = "sync")]
133fn shared_search_pool(num_threads: usize) -> Result<Arc<rayon::ThreadPool>> {
134 if num_threads == 0 {
135 return Err(crate::Error::Internal(
136 "IndexConfig.num_threads must be greater than zero".into(),
137 ));
138 }
139
140 let mut pools = SEARCH_CPU_POOLS
141 .get_or_init(|| parking_lot::Mutex::new(HashMap::new()))
142 .lock();
143 if let Some(pool) = pools.get(&num_threads).and_then(Weak::upgrade) {
144 return Ok(pool);
145 }
146
147 let pool = Arc::new(
151 rayon::ThreadPoolBuilder::new()
152 .num_threads(num_threads)
153 .thread_name(move |idx| format!("hermes-search-{}-{}", num_threads, idx))
154 .build()
155 .map_err(|error| {
156 crate::Error::Internal(format!(
157 "failed to create {num_threads}-thread search pool: {error}"
158 ))
159 })?,
160 );
161 pools.retain(|_, pool| pool.strong_count() > 0);
162 pools.insert(num_threads, Arc::downgrade(&pool));
163 log::info!("[search] process-wide CPU pool: {} thread(s)", num_threads);
164 Ok(pool)
165}
166
167impl Default for IndexConfig {
168 fn default() -> Self {
169 #[cfg(feature = "native")]
170 let compression_threads = crate::default_compression_threads();
171 #[cfg(not(feature = "native"))]
172 let compression_threads = 1;
173
174 #[cfg(feature = "native")]
175 let search_threads = crate::default_search_threads();
176 #[cfg(not(feature = "native"))]
177 let search_threads = 1;
178
179 Self {
180 num_threads: search_threads,
181 num_indexing_threads: 1, num_compression_threads: compression_threads,
183 term_cache_blocks: 256,
184 store_cache_blocks: 32,
185 max_indexing_memory_bytes: 256 * 1024 * 1024, vector_training_max_samples: 10_000_000,
187 #[cfg(target_pointer_width = "64")]
188 vector_training_memory_bytes: 4 * 1024 * 1024 * 1024,
189 #[cfg(not(target_pointer_width = "64"))]
190 vector_training_memory_bytes: usize::MAX,
191 merge_policy: Box::new(crate::merge::TieredMergePolicy::large_scale()),
196 optimization: crate::structures::IndexOptimization::default(),
197 reload_interval_ms: 1000, max_concurrent_merges: 4,
199 #[cfg(feature = "native")]
200 background_merge_permits: Arc::new(tokio::sync::Semaphore::new(4)),
201 merge_bp_time_budget: Some(std::time::Duration::from_secs(600)),
202 #[cfg(target_pointer_width = "64")]
211 bp_memory_budget_bytes: 24 * 1024 * 1024 * 1024,
212 #[cfg(not(target_pointer_width = "64"))]
213 bp_memory_budget_bytes: usize::MAX,
214 #[cfg(feature = "native")]
215 background_reorder_permits: Arc::new(tokio::sync::Semaphore::new(2)),
216 #[cfg(feature = "native")]
217 background_reorder_pool: None,
218 }
219 }
220}
221
222#[cfg(feature = "native")]
231pub struct Index<D: crate::directories::DirectoryWriter + 'static> {
232 directory: Arc<D>,
233 schema: Arc<Schema>,
234 config: IndexConfig,
235 search_resources: searcher::SearcherResources,
237 segment_manager: Arc<crate::merge::SegmentManager<D>>,
239 cached_reader: tokio::sync::OnceCell<IndexReader<D>>,
241}
242
243#[cfg(feature = "native")]
244impl<D: crate::directories::DirectoryWriter + 'static> Index<D> {
245 pub async fn create(directory: D, schema: Schema, config: IndexConfig) -> Result<Self> {
247 let search_resources = searcher::SearcherResources::new(
248 config.term_cache_blocks,
249 config.store_cache_blocks,
250 config.num_threads,
251 )?;
252 let directory = Arc::new(directory);
253 let schema = Arc::new(schema);
254 directory.set_index_label(schema.index_label());
256
257 if directory
261 .exists(std::path::Path::new(INDEX_META_FILENAME))
262 .await?
263 {
264 return Err(crate::Error::Internal(format!(
265 "refusing to create index: {} already exists in this directory; \
266 use Index::open to open the existing index, or delete the \
267 directory first if you really want to start over",
268 INDEX_META_FILENAME
269 )));
270 }
271
272 let metadata = IndexMetadata::new((*schema).clone());
273
274 let segment_manager = Arc::new(crate::merge::SegmentManager::new(
275 Arc::clone(&directory),
276 Arc::clone(&schema),
277 metadata,
278 config.merge_policy.clone_box(),
279 config.term_cache_blocks,
280 config.max_concurrent_merges,
281 Arc::clone(&config.background_merge_permits),
282 config.merge_bp_time_budget,
283 config.bp_memory_budget_bytes,
284 Arc::clone(&config.background_reorder_permits),
285 config.background_reorder_pool.clone(),
286 ));
287
288 segment_manager.update_metadata(|_| {}).await?;
290
291 Ok(Self {
292 directory,
293 schema,
294 config,
295 search_resources,
296 segment_manager,
297 cached_reader: tokio::sync::OnceCell::new(),
298 })
299 }
300
301 pub async fn open(directory: D, config: IndexConfig) -> Result<Self> {
303 let search_resources = searcher::SearcherResources::new(
304 config.term_cache_blocks,
305 config.store_cache_blocks,
306 config.num_threads,
307 )?;
308 let directory = Arc::new(directory);
309
310 let metadata = IndexMetadata::load(directory.as_ref()).await?;
312 let schema = Arc::new(metadata.schema.clone());
313 directory.set_index_label(schema.index_label());
315
316 let segment_manager = Arc::new(crate::merge::SegmentManager::new(
317 Arc::clone(&directory),
318 Arc::clone(&schema),
319 metadata,
320 config.merge_policy.clone_box(),
321 config.term_cache_blocks,
322 config.max_concurrent_merges,
323 Arc::clone(&config.background_merge_permits),
324 config.merge_bp_time_budget,
325 config.bp_memory_budget_bytes,
326 Arc::clone(&config.background_reorder_permits),
327 config.background_reorder_pool.clone(),
328 ));
329
330 segment_manager.try_load_and_publish_trained().await?;
332
333 Ok(Self {
334 directory,
335 schema,
336 config,
337 search_resources,
338 segment_manager,
339 cached_reader: tokio::sync::OnceCell::new(),
340 })
341 }
342
343 pub fn schema(&self) -> &Schema {
345 &self.schema
346 }
347
348 pub fn schema_arc(&self) -> &Arc<Schema> {
350 &self.schema
351 }
352
353 pub fn directory(&self) -> &D {
355 &self.directory
356 }
357
358 pub fn segment_manager(&self) -> &Arc<crate::merge::SegmentManager<D>> {
360 &self.segment_manager
361 }
362
363 pub async fn reader(&self) -> Result<&IndexReader<D>> {
368 self.cached_reader
369 .get_or_try_init(|| async {
370 IndexReader::from_segment_manager_with_resources(
371 Arc::clone(&self.schema),
372 Arc::clone(&self.segment_manager),
373 self.config.reload_interval_ms,
374 self.search_resources.clone(),
375 )
376 .await
377 })
378 .await
379 }
380
381 pub fn config(&self) -> &IndexConfig {
383 &self.config
384 }
385
386 pub async fn segment_readers(&self) -> Result<Vec<Arc<crate::segment::SegmentReader>>> {
388 let reader = self.reader().await?;
389 let searcher = reader.searcher().await?;
390 Ok(searcher.segment_readers().to_vec())
391 }
392
393 pub async fn num_docs(&self) -> Result<u32> {
395 let reader = self.reader().await?;
396 let searcher = reader.searcher().await?;
397 Ok(searcher.num_docs())
398 }
399
400 pub fn default_fields(&self) -> Vec<crate::Field> {
402 if !self.schema.default_fields().is_empty() {
403 self.schema.default_fields().to_vec()
404 } else {
405 self.schema
406 .fields()
407 .filter(|(_, entry)| {
408 entry.indexed && entry.field_type == crate::dsl::FieldType::Text
409 })
410 .map(|(field, _)| field)
411 .collect()
412 }
413 }
414
415 pub fn tokenizers(&self) -> Arc<crate::tokenizer::TokenizerRegistry> {
417 Arc::new(crate::tokenizer::TokenizerRegistry::default())
418 }
419
420 pub fn query_parser(&self) -> crate::dsl::QueryLanguageParser {
422 let default_fields = self.default_fields();
423 let tokenizers = self.tokenizers();
424
425 let query_routers = self.schema.query_routers();
426 if !query_routers.is_empty()
427 && let Ok(router) = crate::dsl::QueryFieldRouter::from_rules(query_routers)
428 {
429 return crate::dsl::QueryLanguageParser::with_router(
430 Arc::clone(&self.schema),
431 default_fields,
432 tokenizers,
433 router,
434 );
435 }
436
437 crate::dsl::QueryLanguageParser::new(Arc::clone(&self.schema), default_fields, tokenizers)
438 }
439
440 pub async fn query(
442 &self,
443 query_str: &str,
444 limit: usize,
445 ) -> Result<crate::query::SearchResponse> {
446 self.query_offset(query_str, limit, 0).await
447 }
448
449 pub async fn query_offset(
451 &self,
452 query_str: &str,
453 limit: usize,
454 offset: usize,
455 ) -> Result<crate::query::SearchResponse> {
456 let parser = self.query_parser();
457 let query = parser
458 .parse(query_str)
459 .map_err(crate::error::Error::Query)?;
460 self.search_offset(query.as_ref(), limit, offset).await
461 }
462
463 pub async fn search(
465 &self,
466 query: &dyn crate::query::Query,
467 limit: usize,
468 ) -> Result<crate::query::SearchResponse> {
469 self.search_offset(query, limit, 0).await
470 }
471
472 pub async fn search_offset(
474 &self,
475 query: &dyn crate::query::Query,
476 limit: usize,
477 offset: usize,
478 ) -> Result<crate::query::SearchResponse> {
479 let reader = self.reader().await?;
480 let searcher = reader.searcher().await?;
481
482 #[cfg(feature = "sync")]
483 let (results, total_seen) = {
484 let runtime_flavor = tokio::runtime::Handle::current().runtime_flavor();
488 if runtime_flavor == tokio::runtime::RuntimeFlavor::MultiThread {
489 tokio::task::block_in_place(|| {
490 searcher.search_with_offset_and_count_sync(query, limit, offset)
491 })?
492 } else {
493 searcher.search_with_offset_and_count_sync(query, limit, offset)?
494 }
495 };
496
497 #[cfg(not(feature = "sync"))]
498 let (results, total_seen) = {
499 searcher
500 .search_with_offset_and_count(query, limit, offset)
501 .await?
502 };
503
504 let total_hits = total_seen;
505 let hits: Vec<crate::query::SearchHit> = results
506 .into_iter()
507 .map(|result| crate::query::SearchHit {
508 address: crate::query::DocAddress::new(result.segment_id, result.doc_id),
509 score: result.score,
510 matched_fields: result.extract_ordinals(),
511 })
512 .collect();
513
514 Ok(crate::query::SearchResponse { hits, total_hits })
515 }
516
517 pub async fn get_document(
519 &self,
520 address: &crate::query::DocAddress,
521 ) -> Result<Option<crate::dsl::Document>> {
522 let reader = self.reader().await?;
523 let searcher = reader.searcher().await?;
524 searcher.get_document(address).await
525 }
526
527 pub async fn get_postings(
529 &self,
530 field: crate::Field,
531 term: &[u8],
532 ) -> Result<
533 Vec<(
534 Arc<crate::segment::SegmentReader>,
535 crate::structures::BlockPostingList,
536 )>,
537 > {
538 let segments = self.segment_readers().await?;
539 let mut results = Vec::new();
540
541 for segment in segments {
542 if let Some(postings) = segment.get_postings(field, term).await? {
543 results.push((segment, postings));
544 }
545 }
546
547 Ok(results)
548 }
549}
550
551#[cfg(feature = "native")]
553impl<D: crate::directories::DirectoryWriter + 'static> Index<D> {
554 pub fn writer(&self) -> writer::IndexWriter<D> {
556 writer::IndexWriter::from_index(self)
557 }
558}
559
560#[cfg(test)]
561mod tests;
562
563