1use std::sync::Arc;
7
8use rustc_hash::FxHashMap;
9
10use crate::directories::Directory;
11use crate::dsl::Schema;
12use crate::error::Result;
13use crate::query::LazyGlobalStats;
14use crate::segment::{SegmentId, SegmentReader, TrainedVectorStructures};
15#[cfg(feature = "native")]
16use crate::segment::{SegmentSnapshot, SegmentTracker};
17
18#[cfg(feature = "native")]
22#[derive(Clone)]
23pub(crate) struct SearcherResources {
24 pub(crate) term_cache_blocks: usize,
25 pub(crate) term_cache_budget_bytes: Option<usize>,
26 pub(crate) store_cache: Arc<crate::segment::SharedStoreCache>,
27 pub(crate) sparse_io_gate: Arc<super::SparseIoGate>,
28 pub(crate) sparse_io_concurrency: usize,
29 #[cfg(feature = "sync")]
30 pub(crate) search_pool: Arc<rayon::ThreadPool>,
31}
32
33#[cfg(feature = "native")]
34impl SearcherResources {
35 pub(crate) fn from_config(config: &super::IndexConfig) -> Result<Self> {
37 Self::new(
38 config.term_cache_blocks,
39 config.term_cache_budget_bytes,
40 config.store_cache_budget_bytes,
41 config.num_threads,
42 config.sparse_io_concurrency,
43 )
44 }
45
46 pub(crate) fn new(
51 term_cache_blocks: usize,
52 term_cache_budget_bytes: Option<usize>,
53 store_cache_budget_bytes: usize,
54 num_threads: usize,
55 sparse_io_concurrency: usize,
56 ) -> Result<Self> {
57 super::validate_term_cache_blocks(term_cache_blocks)?;
58 if num_threads == 0 {
59 return Err(crate::Error::Internal(
60 "IndexConfig.num_threads must be greater than zero".into(),
61 ));
62 }
63 if sparse_io_concurrency == 0 {
64 return Err(crate::Error::Internal(
65 "IndexConfig.sparse_io_concurrency must be greater than zero".into(),
66 ));
67 }
68
69 #[cfg(feature = "sync")]
70 let search_pool = super::shared_search_pool(num_threads)?;
71
72 Ok(Self {
73 term_cache_blocks,
74 term_cache_budget_bytes,
75 store_cache: super::shared_store_cache(store_cache_budget_bytes),
76 sparse_io_gate: super::shared_sparse_io_gate(sparse_io_concurrency),
77 sparse_io_concurrency,
78 #[cfg(feature = "sync")]
79 search_pool,
80 })
81 }
82}
83
84pub struct Searcher<D: Directory + 'static> {
89 #[cfg(feature = "native")]
91 _snapshot: SegmentSnapshot,
92 _phantom: std::marker::PhantomData<D>,
94 segments: Vec<Arc<SegmentReader>>,
96 schema: Arc<Schema>,
98 default_fields: Vec<crate::Field>,
100 tokenizers: Arc<crate::tokenizer::TokenizerRegistry>,
102 trained_vectors: Arc<TrainedVectorStructures>,
104 global_stats: Arc<LazyGlobalStats>,
106 segment_map: FxHashMap<u128, usize>,
108 total_docs: u32,
110 #[cfg(feature = "sync")]
112 search_pool: Arc<rayon::ThreadPool>,
113 #[cfg(feature = "native")]
115 sparse_io_gate: Arc<super::SparseIoGate>,
116 #[cfg(feature = "native")]
117 sparse_io_concurrency: usize,
118}
119
120impl<D: Directory + 'static> Searcher<D> {
121 pub async fn open(
126 directory: Arc<D>,
127 schema: Arc<Schema>,
128 segment_ids: &[String],
129 term_cache_blocks: usize,
130 ) -> Result<Self> {
131 const STANDALONE_STORE_CACHE_BYTES: usize = 32 * 1024 * 1024;
132 #[cfg(feature = "native")]
133 let store_cache = super::shared_store_cache(STANDALONE_STORE_CACHE_BYTES);
134 #[cfg(not(feature = "native"))]
135 let store_cache = Arc::new(crate::segment::SharedStoreCache::new(
136 STANDALONE_STORE_CACHE_BYTES,
137 ));
138 Self::create(
139 directory,
140 schema,
141 segment_ids,
142 Arc::new(TrainedVectorStructures::default()),
143 term_cache_blocks,
144 store_cache,
145 )
146 .await
147 }
148
149 #[cfg(feature = "native")]
151 pub(crate) async fn from_snapshot(
152 directory: Arc<D>,
153 schema: Arc<Schema>,
154 snapshot: SegmentSnapshot,
155 trained_vectors: Arc<TrainedVectorStructures>,
156 resources: SearcherResources,
157 ) -> Result<Self> {
158 let (segments, default_fields, global_stats, segment_map, total_docs) = Self::load_common(
159 &directory,
160 &schema,
161 snapshot.segment_ids(),
162 &trained_vectors,
163 resources.term_cache_blocks,
164 resources.term_cache_budget_bytes,
165 Arc::clone(&resources.store_cache),
166 &[],
167 snapshot.deletions(),
168 )
169 .await?;
170
171 Ok(Self {
172 _snapshot: snapshot,
173 _phantom: std::marker::PhantomData,
174 segments,
175 schema,
176 default_fields,
177 tokenizers: Arc::new(crate::tokenizer::TokenizerRegistry::default()),
178 trained_vectors,
179 global_stats,
180 segment_map,
181 total_docs,
182 #[cfg(feature = "sync")]
183 search_pool: resources.search_pool,
184 sparse_io_gate: resources.sparse_io_gate,
185 sparse_io_concurrency: resources.sparse_io_concurrency,
186 })
187 }
188
189 #[cfg(feature = "native")]
193 pub(crate) async fn from_snapshot_reuse(
194 directory: Arc<D>,
195 schema: Arc<Schema>,
196 snapshot: SegmentSnapshot,
197 trained_vectors: Arc<TrainedVectorStructures>,
198 resources: SearcherResources,
199 existing_segments: &[Arc<SegmentReader>],
200 ) -> Result<Self> {
201 let (segments, default_fields, global_stats, segment_map, total_docs) = Self::load_common(
202 &directory,
203 &schema,
204 snapshot.segment_ids(),
205 &trained_vectors,
206 resources.term_cache_blocks,
207 resources.term_cache_budget_bytes,
208 Arc::clone(&resources.store_cache),
209 existing_segments,
210 snapshot.deletions(),
211 )
212 .await?;
213
214 Ok(Self {
215 _snapshot: snapshot,
216 _phantom: std::marker::PhantomData,
217 segments,
218 schema,
219 default_fields,
220 tokenizers: Arc::new(crate::tokenizer::TokenizerRegistry::default()),
221 trained_vectors,
222 global_stats,
223 segment_map,
224 total_docs,
225 #[cfg(feature = "sync")]
226 search_pool: resources.search_pool,
227 sparse_io_gate: resources.sparse_io_gate,
228 sparse_io_concurrency: resources.sparse_io_concurrency,
229 })
230 }
231
232 async fn create(
234 directory: Arc<D>,
235 schema: Arc<Schema>,
236 segment_ids: &[String],
237 trained_vectors: Arc<TrainedVectorStructures>,
238 term_cache_blocks: usize,
239 store_cache: Arc<crate::segment::SharedStoreCache>,
240 ) -> Result<Self> {
241 let deletions = match super::IndexMetadata::load(directory.as_ref()).await {
242 Ok(metadata) => {
243 if segment_ids.iter().any(|id| !metadata.has_segment(id)) {
244 return Err(crate::Error::Query(
245 "segment list is stale; reopen the index metadata before searching".into(),
246 ));
247 }
248 metadata
249 .segment_metas
250 .into_iter()
251 .filter_map(|(id, info)| info.deletions.map(|d| (id, (info.num_docs, d))))
252 .collect()
253 }
254 Err(crate::Error::Io(error)) if error.kind() == std::io::ErrorKind::NotFound => {
255 Default::default()
256 }
257 Err(error) => return Err(error),
258 };
259 let (segments, default_fields, global_stats, segment_map, total_docs) = Self::load_common(
260 &directory,
261 &schema,
262 segment_ids,
263 &trained_vectors,
264 term_cache_blocks,
265 None,
266 store_cache,
267 &[],
268 &deletions,
269 )
270 .await?;
271
272 #[cfg(feature = "native")]
273 let _snapshot = {
274 let tracker = Arc::new(SegmentTracker::new());
275 SegmentSnapshot::new(tracker, segment_ids.to_vec())
276 };
277
278 #[cfg(feature = "sync")]
279 let search_pool = super::shared_search_pool(crate::default_search_threads())?;
280 #[cfg(feature = "native")]
281 let sparse_io_concurrency = 4;
282 #[cfg(feature = "native")]
283 let sparse_io_gate = super::shared_sparse_io_gate(sparse_io_concurrency);
284
285 let _ = directory; Ok(Self {
287 #[cfg(feature = "native")]
288 _snapshot,
289 _phantom: std::marker::PhantomData,
290 segments,
291 schema,
292 default_fields,
293 tokenizers: Arc::new(crate::tokenizer::TokenizerRegistry::default()),
294 trained_vectors,
295 global_stats,
296 segment_map,
297 total_docs,
298 #[cfg(feature = "sync")]
299 search_pool,
300 #[cfg(feature = "native")]
301 sparse_io_gate,
302 #[cfg(feature = "native")]
303 sparse_io_concurrency,
304 })
305 }
306
307 #[allow(clippy::too_many_arguments)]
309 async fn load_common(
310 directory: &Arc<D>,
311 schema: &Arc<Schema>,
312 segment_ids: &[String],
313 trained_vectors: &Arc<TrainedVectorStructures>,
314 term_cache_blocks: usize,
315 term_cache_budget_bytes: Option<usize>,
316 store_cache: Arc<crate::segment::SharedStoreCache>,
317 existing_segments: &[Arc<SegmentReader>],
318 deletions: &std::collections::HashMap<String, (u32, crate::segment::DeletionMeta)>,
319 ) -> Result<(
320 Vec<Arc<SegmentReader>>,
321 Vec<crate::Field>,
322 Arc<LazyGlobalStats>,
323 FxHashMap<u128, usize>,
324 u32,
325 )> {
326 let segments = Self::load_segments(
327 directory,
328 schema,
329 segment_ids,
330 trained_vectors,
331 term_cache_blocks,
332 term_cache_budget_bytes,
333 store_cache,
334 existing_segments,
335 deletions,
336 )
337 .await?;
338 let default_fields = Self::build_default_fields(schema);
339 let global_stats = Arc::new(LazyGlobalStats::new(segments.clone()));
340 let (segment_map, total_docs) = Self::build_lookup_tables(&segments);
341 Ok((
342 segments,
343 default_fields,
344 global_stats,
345 segment_map,
346 total_docs,
347 ))
348 }
349
350 #[allow(clippy::too_many_arguments)]
354 async fn load_segments(
355 directory: &Arc<D>,
356 schema: &Arc<Schema>,
357 segment_ids: &[String],
358 trained_vectors: &Arc<TrainedVectorStructures>,
359 term_cache_blocks: usize,
360 term_cache_budget_bytes: Option<usize>,
361 store_cache: Arc<crate::segment::SharedStoreCache>,
362 existing_segments: &[Arc<SegmentReader>],
363 deletions: &std::collections::HashMap<String, (u32, crate::segment::DeletionMeta)>,
364 ) -> Result<Vec<Arc<SegmentReader>>> {
365 let existing_map: FxHashMap<u128, Arc<SegmentReader>> = existing_segments
367 .iter()
368 .map(|seg| (seg.meta().id, Arc::clone(seg)))
369 .collect();
370
371 let mut valid_segments: Vec<(usize, SegmentId)> = Vec::with_capacity(segment_ids.len());
376 for (idx, id_str) in segment_ids.iter().enumerate() {
377 let sid = SegmentId::from_hex(id_str).ok_or_else(|| {
378 crate::error::Error::Corruption(format!(
379 "Invalid segment ID in metadata: {id_str:?}"
380 ))
381 })?;
382 valid_segments.push((idx, sid));
383 }
384
385 let mut loaded = Vec::with_capacity(valid_segments.len());
387 let mut to_load = Vec::new();
388 let mut refreshed = 0;
389 for (idx, sid) in valid_segments {
390 let deletion = deletions.get(&sid.to_hex());
391 let existing = existing_map.get(&sid.0);
392 if let (Some(existing), Some((num_docs, _))) = (existing, deletion)
393 && *num_docs != existing.num_docs()
394 {
395 return Err(crate::Error::Corruption(
396 "deletion metadata row count mismatch".into(),
397 ));
398 }
399 if let Some(existing) = existing
400 && existing.deletion_meta() == deletion.map(|(_, meta)| meta)
401 {
402 loaded.push((idx, Arc::clone(existing)));
403 } else {
404 refreshed += usize::from(existing.is_some());
405 to_load.push((idx, sid, existing.cloned()));
406 }
407 }
408
409 if !existing_segments.is_empty() {
410 log::info!(
411 "[searcher] index={} reusing {} segment readers, refreshing {} visibility views, loading {} new",
412 schema.index_label(),
413 loaded.len(),
414 refreshed,
415 to_load.len() - refreshed,
416 );
417 }
418
419 const MAX_CONCURRENT_SEGMENT_OPENS: usize = 2;
422 use futures::{StreamExt, TryStreamExt};
423 let store_cache_directory_namespace = Arc::as_ptr(directory) as usize;
425 let results: Vec<_> =
426 futures::stream::iter(to_load.into_iter().map(|(idx, sid, existing)| {
427 let store_cache = Arc::clone(&store_cache);
428 async move {
429 let deletion = deletions.get(&sid.to_hex());
430 let mut reader = match existing {
431 Some(existing) => {
432 existing
433 .with_deletions(
434 directory.as_ref(),
435 deletion.map(|(_, meta)| meta.clone()),
436 )
437 .await?
438 }
439 None => {
440 let mut reader = SegmentReader::open_with_store_cache(
441 directory.as_ref(),
442 sid,
443 Arc::clone(schema),
444 term_cache_blocks,
445 term_cache_budget_bytes,
446 store_cache_directory_namespace,
447 store_cache,
448 )
449 .await
450 .map_err(|error| {
451 crate::Error::Internal(format!(
452 "Failed to open segment {:016x}: {:?}",
453 sid.0, error
454 ))
455 })?;
456 if let Some((num_docs, meta)) = deletion {
457 if *num_docs != reader.num_docs() {
458 return Err(crate::Error::Corruption(
459 "deletion metadata row count mismatch".into(),
460 ));
461 }
462 reader
463 .load_deletions(directory.as_ref(), meta.clone())
464 .await?;
465 }
466 reader
467 }
468 };
469 reader.set_trained_vectors(Arc::clone(trained_vectors));
470 Ok((idx, Arc::new(reader)))
471 }
472 }))
473 .buffer_unordered(MAX_CONCURRENT_SEGMENT_OPENS)
474 .try_collect()
475 .await?;
476 loaded.extend(results);
477
478 loaded.sort_by_key(|(idx, _)| *idx);
480
481 let segments: Vec<Arc<SegmentReader>> = loaded.into_iter().map(|(_, seg)| seg).collect();
482
483 let total_docs: u64 = segments.iter().map(|s| s.meta().num_docs as u64).sum();
487 let mut total_heap = 0usize;
488 let mut total_file_backed = 0u64;
489 let mut total_pinned = 0u64;
490 let mut total_pin_intended = 0u64;
491 for seg in &segments {
492 let stats = seg.memory_stats();
493 let heap = stats.estimated_heap_bytes();
494 let file_backed = stats.file_backed_bytes();
495 total_heap = total_heap.saturating_add(heap);
496 total_file_backed = total_file_backed.saturating_add(file_backed);
497 total_pinned = total_pinned.saturating_add(stats.pinned_metadata_bytes);
498 total_pin_intended = total_pin_intended.saturating_add(stats.pin_intended_bytes);
499 log::info!(
500 "[searcher] index={} segment {:016x}: docs={}, heap_estimate={} \
501 (term_cache={}, store_cache={}, sparse_vectors={}, dense_vectors={}), \
502 file_backed={} (term_bloom={}, sparse_vectors={}, dense_vectors={}), \
503 pinned_metadata={} of {} eligible \
504 (sparse_vectors={} of {}, dense_vectors={} of {})",
505 schema.index_label(),
506 stats.segment_id,
507 stats.num_docs,
508 crate::format_bytes(heap as u64),
509 crate::format_bytes(stats.term_dict_cache_bytes as u64),
510 crate::format_bytes(stats.store_cache_bytes as u64),
511 crate::format_bytes(stats.sparse_heap_bytes as u64),
512 crate::format_bytes(stats.dense_heap_bytes as u64),
513 crate::format_bytes(file_backed),
514 crate::format_bytes(stats.term_bloom_file_bytes),
515 crate::format_bytes(stats.sparse_file_backed_bytes),
516 crate::format_bytes(stats.dense_file_backed_bytes),
517 crate::format_bytes(stats.pinned_metadata_bytes),
518 crate::format_bytes(stats.pin_intended_bytes),
519 crate::format_bytes(stats.sparse_pinned_metadata_bytes),
520 crate::format_bytes(stats.sparse_pin_intended_bytes),
521 crate::format_bytes(stats.dense_pinned_metadata_bytes),
522 crate::format_bytes(stats.dense_pin_intended_bytes),
523 );
524 }
525 let rss_bytes = process_rss_bytes();
527 log::info!(
528 "[searcher] index={} loaded {} segments: total_docs={}, heap_estimate={}, \
529 file_backed={}, pinned_metadata={} of {} eligible, \
530 shared_store_cache={} in {} blocks, process_rss={}",
531 schema.index_label(),
532 segments.len(),
533 total_docs,
534 crate::format_bytes(total_heap as u64),
535 crate::format_bytes(total_file_backed),
536 crate::format_bytes(total_pinned),
537 crate::format_bytes(total_pin_intended),
538 crate::format_bytes(store_cache.total_bytes() as u64),
539 store_cache.total_blocks(),
540 crate::format_bytes(rss_bytes),
541 );
542
543 let mut ann_per_field: std::collections::BTreeMap<u32, (u64, u64, u32, u32, u64)> =
547 std::collections::BTreeMap::new();
548 for segment in segments.iter() {
549 for &field_id in segment.vector_indexes().keys() {
550 if let Some(health) = segment.ann_health(crate::Field(field_id)) {
551 let entry = ann_per_field.entry(field_id).or_default();
552 entry.0 += health.vectors;
553 entry.1 += health.payload_bytes;
554 entry.2 += health.runs;
555 entry.3 += health.clusters_nonempty;
556 entry.4 = entry.4.max(health.largest_cluster_vectors);
557 }
558 }
559 }
560 for (field_id, (vectors, payload, runs, clusters, largest)) in ann_per_field {
561 log::info!(
562 "[ann_health] index={} field={field_id} aggregate: vectors={vectors} \
563 payload={} runs={runs} fragmentation={:.2} worst_leaf_vectors={largest}",
564 schema.index_label(),
565 crate::format_bytes(payload),
566 if clusters == 0 {
567 0.0
568 } else {
569 f64::from(runs) / f64::from(clusters)
570 },
571 );
572 }
573
574 Ok(segments)
575 }
576
577 fn build_default_fields(schema: &Schema) -> Vec<crate::Field> {
579 if !schema.default_fields().is_empty() {
580 schema.default_fields().to_vec()
581 } else {
582 schema
583 .fields()
584 .filter(|(_, entry)| {
585 entry.indexed && entry.field_type == crate::dsl::FieldType::Text
586 })
587 .map(|(field, _)| field)
588 .collect()
589 }
590 }
591
592 pub fn schema(&self) -> &Schema {
594 &self.schema
595 }
596
597 pub fn schema_arc(&self) -> Arc<Schema> {
598 Arc::clone(&self.schema)
599 }
600
601 pub fn segment_readers(&self) -> &[Arc<SegmentReader>] {
603 &self.segments
604 }
605
606 pub fn default_fields(&self) -> &[crate::Field] {
608 &self.default_fields
609 }
610
611 pub fn tokenizers(&self) -> &crate::tokenizer::TokenizerRegistry {
613 &self.tokenizers
614 }
615
616 pub fn trained_centroids(&self) -> &FxHashMap<u32, Arc<crate::structures::CoarseCentroids>> {
618 &self.trained_vectors.centroids
619 }
620
621 pub fn trained_binary_quantizers(
622 &self,
623 ) -> &FxHashMap<u32, Arc<crate::structures::BinaryCoarseQuantizer>> {
624 &self.trained_vectors.binary_quantizers
625 }
626
627 pub fn global_stats(&self) -> &Arc<LazyGlobalStats> {
629 &self.global_stats
630 }
631
632 fn build_lookup_tables(segments: &[Arc<SegmentReader>]) -> (FxHashMap<u128, usize>, u32) {
634 let mut segment_map = FxHashMap::default();
635 let mut total = 0u32;
636 for (i, seg) in segments.iter().enumerate() {
637 segment_map.insert(seg.meta().id, i);
638 total = total.saturating_add(seg.num_live_docs());
639 }
640 (segment_map, total)
641 }
642
643 pub fn num_docs(&self) -> u32 {
645 self.total_docs
646 }
647
648 pub fn segment_map(&self) -> &FxHashMap<u128, usize> {
650 &self.segment_map
651 }
652
653 #[cfg(feature = "sync")]
655 pub(crate) fn install_search_cpu<R: Send>(&self, operation: impl FnOnce() -> R + Send) -> R {
656 #[cfg(not(feature = "query-diagnostics"))]
657 {
658 self.search_pool.install(operation)
659 }
660 #[cfg(feature = "query-diagnostics")]
661 {
662 let submitted = std::time::Instant::now();
663 let (result, queued, finished) = self.search_pool.install(|| {
664 let queued = submitted.elapsed();
665 let result = operation();
666 (result, queued, std::time::Instant::now())
667 });
668 let returned = finished.elapsed();
669 crate::observe::search_work!(search_pool_installs += 1);
670 crate::observe::search_work!(
671 search_pool_queue_ns += queued.as_nanos().min(u128::from(u64::MAX)) as u64
672 );
673 crate::observe::search_work!(
674 search_pool_return_ns += returned.as_nanos().min(u128::from(u64::MAX)) as u64
675 );
676 result
677 }
678 }
679
680 #[cfg(not(feature = "sync"))]
684 pub(crate) fn install_search_cpu<R>(&self, operation: impl FnOnce() -> R) -> R {
685 operation()
686 }
687
688 #[cfg(feature = "sync")]
692 pub(crate) async fn run_search_cpu<F>(&self, future: F) -> F::Output
693 where
694 F: std::future::Future + Send,
695 F::Output: Send,
696 {
697 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
698 return future.await;
699 };
700 if runtime.runtime_flavor() != tokio::runtime::RuntimeFlavor::MultiThread {
701 return future.await;
702 }
703 let mut future = std::pin::pin!(future);
704 futures::future::poll_fn(|context| {
705 let waker = context.waker().clone();
706 tokio::task::block_in_place(|| {
707 self.install_search_cpu(|| {
708 let _entered = runtime.enter();
709 future
710 .as_mut()
711 .poll(&mut std::task::Context::from_waker(&waker))
712 })
713 })
714 })
715 .await
716 }
717
718 #[cfg(not(feature = "sync"))]
719 pub(crate) async fn run_search_cpu<F: std::future::Future>(&self, future: F) -> F::Output {
720 future.await
721 }
722
723 pub fn num_segments(&self) -> usize {
725 self.segments.len()
726 }
727
728 pub async fn doc(&self, segment_id: u128, doc_id: u32) -> Result<Option<crate::dsl::Document>> {
730 if let Some(&idx) = self.segment_map.get(&segment_id) {
731 return self.segments[idx].doc(doc_id).await;
732 }
733 Ok(None)
734 }
735
736 pub async fn search(
738 &self,
739 query: &dyn crate::query::Query,
740 limit: usize,
741 ) -> Result<Vec<crate::query::SearchResult>> {
742 let (results, _) = self.search_with_count(query, limit).await?;
743 Ok(results)
744 }
745
746 pub async fn search_with_count(
749 &self,
750 query: &dyn crate::query::Query,
751 limit: usize,
752 ) -> Result<(Vec<crate::query::SearchResult>, u32)> {
753 self.search_with_offset_and_count(query, limit, 0).await
754 }
755
756 pub async fn search_with_offset(
758 &self,
759 query: &dyn crate::query::Query,
760 limit: usize,
761 offset: usize,
762 ) -> Result<Vec<crate::query::SearchResult>> {
763 let (results, _) = self
764 .search_with_offset_and_count(query, limit, offset)
765 .await?;
766 Ok(results)
767 }
768
769 pub async fn search_with_offset_and_count(
771 &self,
772 query: &dyn crate::query::Query,
773 limit: usize,
774 offset: usize,
775 ) -> Result<(Vec<crate::query::SearchResult>, u32)> {
776 self.search_internal(query, limit, offset, false).await
777 }
778
779 pub async fn search_with_positions(
783 &self,
784 query: &dyn crate::query::Query,
785 limit: usize,
786 ) -> Result<(Vec<crate::query::SearchResult>, u32)> {
787 self.search_internal(query, limit, 0, true).await
788 }
789
790 pub async fn search_with_positions_budgeted(
795 &self,
796 query: &dyn crate::query::Query,
797 limit: usize,
798 deadline: Option<std::time::Instant>,
799 ) -> Result<(Vec<crate::query::SearchResult>, u32, bool)> {
800 self.search_internal_budgeted(query, limit, 0, true, deadline, None)
801 .await
802 }
803
804 pub async fn search_with_count_budgeted(
807 &self,
808 query: &dyn crate::query::Query,
809 limit: usize,
810 deadline: Option<std::time::Instant>,
811 ) -> Result<(Vec<crate::query::SearchResult>, u32, bool)> {
812 self.search_internal_budgeted(query, limit, 0, false, deadline, None)
813 .await
814 }
815
816 pub async fn search_with_positions_budgeted_stats(
820 &self,
821 query: &dyn crate::query::Query,
822 limit: usize,
823 deadline: Option<std::time::Instant>,
824 stats: Option<Arc<crate::query::GlobalStats>>,
825 ) -> Result<(Vec<crate::query::SearchResult>, u32, bool)> {
826 self.search_internal_budgeted(query, limit, 0, true, deadline, stats)
827 .await
828 }
829
830 pub async fn search_with_count_budgeted_stats(
833 &self,
834 query: &dyn crate::query::Query,
835 limit: usize,
836 deadline: Option<std::time::Instant>,
837 stats: Option<Arc<crate::query::GlobalStats>>,
838 ) -> Result<(Vec<crate::query::SearchResult>, u32, bool)> {
839 self.search_internal_budgeted(query, limit, 0, false, deadline, stats)
840 .await
841 }
842
843 pub fn query_text_stats(
848 &self,
849 query: &dyn crate::query::Query,
850 stats_override: Option<Arc<crate::query::GlobalStats>>,
851 ) -> Option<Arc<crate::query::GlobalStats>> {
852 if stats_override.is_some() {
853 return stats_override;
854 }
855 if self.segments.len() < 2 {
856 return None;
857 }
858 let mut terms = Vec::new();
859 query.text_terms(&mut terms);
860 if terms.is_empty() {
861 return None;
862 }
863 terms.sort_unstable_by(|a, b| (a.0.0, &a.1).cmp(&(b.0.0, &b.1)));
864 terms.dedup_by(|a, b| a.0.0 == b.0.0 && a.1 == b.1);
865 Some(Arc::new(self.global_stats.text_stats_for(&terms)))
866 }
867
868 fn prepare_global_lsp(
875 &self,
876 query: &dyn crate::query::Query,
877 retrieval_depth: usize,
878 parallel: bool,
879 ) -> Result<Vec<Option<std::sync::Arc<crate::query::bmp::LspSegmentPlan>>>> {
880 let total_start = crate::observe::WallTimer::start();
881 let empty = || vec![None; self.segments.len()];
882 if retrieval_depth == 0 {
883 return Ok(empty());
884 }
885 let crate::query::QueryDecomposition::SparseTerms(infos) = query.sparse_decomposition()
886 else {
887 return Ok(empty());
888 };
889 let Some(&first) = infos.first() else {
890 return Ok(empty());
891 };
892 if infos
893 .iter()
894 .any(|info| info.field != first.field || info.lsp_gamma != first.lsp_gamma)
895 {
896 return Ok(empty());
897 }
898 let field = first.field;
899 let field_label = self.schema.get_field_name(field).unwrap_or("?");
900 let (total_superblocks, total_coarse_groups, planning_depth) = self
901 .segments
902 .iter()
903 .filter_map(|segment| segment.bmp_index(field))
904 .fold(
905 (0usize, 0usize, retrieval_depth),
906 |(total, coarse, depth), bmp| {
907 (
908 total.saturating_add(bmp.num_superblocks as usize),
909 coarse.saturating_add(bmp.num_coarse_groups as usize),
910 depth.max(crate::query::bmp_executor_limit(
911 retrieval_depth,
912 first.over_fetch_factor,
913 bmp,
914 )),
915 )
916 },
917 );
918 let Some(reference_bmp) = self
919 .segments
920 .iter()
921 .find_map(|segment| segment.bmp_index(field))
922 else {
923 return Ok(empty());
924 };
925 if !infos.iter().any(|info| info.candidate) {
926 return Ok(empty());
927 }
928 let prepare_start = crate::observe::WallTimer::start();
929 let Some(prepared_query) =
930 crate::query::bmp::prepare_bmp_query_infos(reference_bmp.dims(), &infos)?
931 else {
932 return Ok(empty());
933 };
934 let infos: std::sync::Arc<[crate::query::SparseTermQueryInfo]> = infos.into();
935 let prepared_query = std::sync::Arc::new(prepared_query);
936 let prepare_secs = prepare_start.secs();
937
938 let local_plans = || {
939 let plan = std::sync::Arc::new(crate::query::bmp::LspSegmentPlan {
940 infos: std::sync::Arc::clone(&infos),
941 prepared_query: std::sync::Arc::clone(&prepared_query),
942 selection: None,
943 });
944 self.segments
945 .iter()
946 .map(|segment| {
947 segment
948 .bmp_index(field)
949 .map(|_| std::sync::Arc::clone(&plan))
950 })
951 .collect()
952 };
953 let gamma = first
954 .lsp_gamma
955 .unwrap_or_else(|| crate::query::bmp::recommended_lsp_gamma(planning_depth));
956 if gamma == 0 || gamma >= total_superblocks {
957 crate::observe::bmp_lsp(
961 self.schema.index_label(),
962 field_label,
963 total_start.secs(),
964 prepare_secs,
965 0.0,
966 0.0,
967 total_superblocks,
968 gamma,
969 total_coarse_groups,
970 0,
971 0,
972 );
973 return Ok(local_plans());
974 }
975 let hierarchy_scan_start = crate::observe::WallTimer::start();
976 let prepare = |segment: &std::sync::Arc<crate::segment::SegmentReader>| {
977 segment
978 .bmp_index(field)
979 .map(|bmp| crate::query::bmp::prepare_lsp_coarse_ubs(bmp, &prepared_query))
980 .transpose()
981 };
982
983 #[cfg(feature = "sync")]
984 let coarse_bounds: Vec<Option<Vec<f32>>> = if parallel {
985 use rayon::prelude::*;
986 self.search_pool.install(|| {
987 self.segments
988 .par_iter()
989 .map(prepare)
990 .collect::<Result<Vec<_>>>()
991 })?
992 } else {
993 self.segments
994 .iter()
995 .map(prepare)
996 .collect::<Result<Vec<_>>>()?
997 };
998 #[cfg(not(feature = "sync"))]
999 let coarse_bounds: Vec<Option<Vec<f32>>> = {
1000 let _ = parallel;
1001 self.segments
1002 .iter()
1003 .map(prepare)
1004 .collect::<Result<Vec<_>>>()?
1005 };
1006 let hierarchy_scan_secs = hierarchy_scan_start.secs();
1007
1008 let select_start = crate::observe::WallTimer::start();
1009 let selection =
1010 select_global_lsp_hierarchical(&coarse_bounds, gamma, |segment, group, out| {
1011 let bmp = self.segments[segment].bmp_index(field).ok_or_else(|| {
1012 crate::Error::Internal(
1013 "BMP coarse plan references a segment without the sparse field".into(),
1014 )
1015 })?;
1016 crate::query::bmp::expand_lsp_coarse_group(bmp, &prepared_query, group, out)
1017 })?;
1018 let mut plans = Vec::with_capacity(self.segments.len());
1019 for (segment, selected) in selection.selected.into_iter().enumerate() {
1020 if self.segments[segment].bmp_index(field).is_none() {
1021 plans.push(None);
1022 continue;
1023 }
1024 let (selected_superblocks, selected_bounds): (Vec<_>, Vec<_>) =
1025 selected.into_iter().unzip();
1026 plans.push(Some(std::sync::Arc::new(
1027 crate::query::bmp::LspSegmentPlan {
1028 infos: std::sync::Arc::clone(&infos),
1029 prepared_query: std::sync::Arc::clone(&prepared_query),
1030 selection: Some(crate::query::bmp::LspSelection {
1031 sb_ubs: selected_bounds,
1032 sb_order: selected_superblocks,
1033 }),
1034 },
1035 )));
1036 }
1037 let select_secs = select_start.secs();
1038 crate::observe::bmp_lsp(
1039 self.schema.index_label(),
1040 field_label,
1041 total_start.secs(),
1042 prepare_secs,
1043 hierarchy_scan_secs,
1044 select_secs,
1045 total_superblocks,
1046 gamma,
1047 total_coarse_groups,
1048 selection.expanded_groups,
1049 selection.evaluated_superblocks,
1050 );
1051 log::debug!(
1052 "[searcher] BMP hierarchical LSP: index={}, field={}, coarse_groups={}/{}, E_superblocks={}/{}, gamma={}",
1053 self.schema.index_label(),
1054 field_label,
1055 selection.expanded_groups,
1056 coarse_bounds
1057 .iter()
1058 .filter_map(Option::as_ref)
1059 .map(Vec::len)
1060 .sum::<usize>(),
1061 selection.evaluated_superblocks,
1062 total_superblocks,
1063 gamma,
1064 );
1065 Ok(plans)
1066 }
1067
1068 fn ordered_lsp_segments(
1069 &self,
1070 plans: &[Option<std::sync::Arc<crate::query::bmp::LspSegmentPlan>>],
1071 ) -> Vec<usize> {
1072 let mut order: Vec<usize> = (0..self.segments.len()).collect();
1073 order.sort_unstable_by(|&left, &right| {
1074 let left_priority = plans[left]
1075 .as_ref()
1076 .map_or(f32::NEG_INFINITY, |plan| plan.priority());
1077 let right_priority = plans[right]
1078 .as_ref()
1079 .map_or(f32::NEG_INFINITY, |plan| plan.priority());
1080 right_priority
1081 .total_cmp(&left_priority)
1082 .then_with(|| {
1083 self.segments[right]
1084 .num_docs()
1085 .cmp(&self.segments[left].num_docs())
1086 })
1087 .then_with(|| {
1088 self.segments[left]
1089 .meta()
1090 .id
1091 .cmp(&self.segments[right].meta().id)
1092 })
1093 });
1094 order
1095 }
1096
1097 #[inline]
1098 fn bmp_wave_width(&self) -> usize {
1099 #[cfg(feature = "native")]
1100 {
1101 self.sparse_io_concurrency
1102 }
1103 #[cfg(not(feature = "native"))]
1104 {
1105 4
1106 }
1107 }
1108
1109 async fn search_internal(
1111 &self,
1112 query: &dyn crate::query::Query,
1113 limit: usize,
1114 offset: usize,
1115 collect_positions: bool,
1116 ) -> Result<(Vec<crate::query::SearchResult>, u32)> {
1117 let (results, seen, _) = self
1118 .search_internal_budgeted(query, limit, offset, collect_positions, None, None)
1119 .await?;
1120 Ok((results, seen))
1121 }
1122
1123 async fn search_internal_budgeted(
1124 &self,
1125 query: &dyn crate::query::Query,
1126 limit: usize,
1127 offset: usize,
1128 collect_positions: bool,
1129 deadline: Option<std::time::Instant>,
1130 stats_override: Option<Arc<crate::query::GlobalStats>>,
1131 ) -> Result<(Vec<crate::query::SearchResult>, u32, bool)> {
1132 let fetch_limit = checked_search_window(limit, offset)?;
1133 let text_stats = self.query_text_stats(query, stats_override);
1134
1135 #[cfg(feature = "sync")]
1141 if !self.segments.is_empty()
1142 && tokio::runtime::Handle::current().runtime_flavor()
1143 == tokio::runtime::RuntimeFlavor::MultiThread
1144 {
1145 return self.search_internal_parallel(
1146 query,
1147 fetch_limit,
1148 offset,
1149 collect_positions,
1150 deadline,
1151 text_stats,
1152 );
1153 }
1154
1155 const MAX_ASYNC_SEGMENT_SEARCHES: usize = 8;
1159 use futures::StreamExt;
1160 use futures::TryStreamExt;
1161 let shared = crate::query::SharedThreshold::for_limit(fetch_limit).with_deadline(deadline);
1166 let lsp_plans = self.prepare_global_lsp(query, fetch_limit, false)?;
1167 let mut total_seen: u32 = 0;
1168 let mut merged = Vec::new();
1169 let mut merge_scratch = Vec::new();
1170 let bmp_planned = lsp_plans.iter().any(Option::is_some);
1171 let order = if bmp_planned {
1172 self.ordered_lsp_segments(&lsp_plans)
1173 } else {
1174 (0..self.segments.len()).collect()
1175 };
1176 #[cfg(feature = "native")]
1177 let sparse_io_gate = Arc::clone(&self.sparse_io_gate);
1178 let run_segment = |segment_index: usize| {
1179 let text_stats = text_stats.clone();
1180 let segment = Arc::clone(&self.segments[segment_index]);
1181 let lsp_plan = lsp_plans[segment_index].clone();
1182 let shared = shared.clone();
1183 #[cfg(feature = "native")]
1184 let sparse_io_gate = Arc::clone(&sparse_io_gate);
1185 async move {
1186 if lsp_plan.as_ref().is_some_and(|plan| !plan.has_work()) {
1187 return Ok((Vec::new(), 0u32));
1188 }
1189 #[cfg(feature = "native")]
1190 let _io_permit = if lsp_plan.is_some() {
1191 Some(sparse_io_gate.acquire_async().await)
1192 } else {
1193 None
1194 };
1195 let sid = segment.meta().id;
1196 let (mut results, segment_seen) = crate::query::search_segment_shared_planned(
1197 segment.as_ref(),
1198 query,
1199 fetch_limit,
1200 collect_positions,
1201 shared.clone(),
1202 lsp_plan,
1203 text_stats.clone(),
1204 )
1205 .await?;
1206 if fetch_limit > 0 && results.len() >= fetch_limit {
1207 shared.raise(results[fetch_limit - 1].score);
1208 }
1209 for result in &mut results {
1210 result.segment_id = sid;
1211 }
1212 Ok::<_, crate::error::Error>((results, segment_seen))
1213 }
1214 };
1215
1216 let mut remainder = order.as_slice();
1217 if bmp_planned {
1218 if let Some((&pilot, rest)) = order.split_first() {
1221 let (batch, segment_seen) = run_segment(pilot).await?;
1222 total_seen = total_seen.saturating_add(segment_seen);
1223 merge_ranked_reuse(&mut merged, batch, fetch_limit, &mut merge_scratch);
1224 remainder = rest;
1225 }
1226 }
1227 let concurrency = if bmp_planned {
1228 self.bmp_wave_width()
1229 } else {
1230 MAX_ASYNC_SEGMENT_SEARCHES
1231 };
1232 let searches = futures::stream::iter(remainder.iter().copied().map(run_segment))
1233 .buffer_unordered(concurrency);
1234 futures::pin_mut!(searches);
1235 while let Some((batch, segment_seen)) = searches.try_next().await? {
1236 total_seen = total_seen.saturating_add(segment_seen);
1237 merge_ranked_reuse(&mut merged, batch, fetch_limit, &mut merge_scratch);
1238 }
1239
1240 let results = apply_result_offset(merged, fetch_limit, offset);
1241 Ok((results, total_seen, shared.truncated()))
1242 }
1243
1244 #[cfg(feature = "sync")]
1249 fn search_internal_parallel(
1250 &self,
1251 query: &dyn crate::query::Query,
1252 fetch_limit: usize,
1253 offset: usize,
1254 collect_positions: bool,
1255 deadline: Option<std::time::Instant>,
1256 text_stats: Option<Arc<crate::query::GlobalStats>>,
1257 ) -> Result<(Vec<crate::query::SearchResult>, u32, bool)> {
1258 tokio::task::block_in_place(|| {
1259 self.search_internal_sync_budgeted(
1260 query,
1261 fetch_limit,
1262 offset,
1263 collect_positions,
1264 deadline,
1265 text_stats,
1266 )
1267 })
1268 }
1269
1270 #[cfg(feature = "sync")]
1275 fn search_internal_sync(
1276 &self,
1277 query: &dyn crate::query::Query,
1278 fetch_limit: usize,
1279 offset: usize,
1280 collect_positions: bool,
1281 ) -> Result<(Vec<crate::query::SearchResult>, u32)> {
1282 let text_stats = self.query_text_stats(query, None);
1283 let (results, total_seen, _) = self.search_internal_sync_budgeted(
1284 query,
1285 fetch_limit,
1286 offset,
1287 collect_positions,
1288 None,
1289 text_stats,
1290 )?;
1291 Ok((results, total_seen))
1292 }
1293
1294 #[cfg(feature = "sync")]
1295 fn search_internal_sync_budgeted(
1296 &self,
1297 query: &dyn crate::query::Query,
1298 fetch_limit: usize,
1299 offset: usize,
1300 collect_positions: bool,
1301 deadline: Option<std::time::Instant>,
1302 text_stats: Option<Arc<crate::query::GlobalStats>>,
1303 ) -> Result<(Vec<crate::query::SearchResult>, u32, bool)> {
1304 let (merged, total_seen, truncated) =
1305 self.search_segments_sync(query, fetch_limit, collect_positions, deadline, text_stats)?;
1306 let results = apply_result_offset(merged, fetch_limit, offset);
1307 Ok((results, total_seen, truncated))
1308 }
1309
1310 #[cfg(feature = "sync")]
1318 fn search_segments_sync(
1319 &self,
1320 query: &dyn crate::query::Query,
1321 fetch_limit: usize,
1322 collect_positions: bool,
1323 deadline: Option<std::time::Instant>,
1324 text_stats: Option<Arc<crate::query::GlobalStats>>,
1325 ) -> Result<(Vec<crate::query::SearchResult>, u32, bool)> {
1326 use rayon::prelude::*;
1327
1328 let lsp_plans = self.prepare_global_lsp(query, fetch_limit, true)?;
1329 let shared = crate::query::SharedThreshold::for_limit(fetch_limit).with_deadline(deadline);
1330 #[cfg(feature = "query-diagnostics")]
1331 let diagnostic_context = crate::search_diagnostics::WorkContext::current();
1332 let run_segment = |segment_index: &usize| {
1333 #[cfg(feature = "query-diagnostics")]
1334 let _diagnostic_scope = diagnostic_context.as_ref().map(|context| context.enter());
1335 let segment = &self.segments[*segment_index];
1336 let lsp_plan = lsp_plans[*segment_index].clone();
1337 if lsp_plan.as_ref().is_some_and(|plan| !plan.has_work()) {
1338 return Ok((Vec::new(), 0u32));
1339 }
1340 let _io_permit = lsp_plan.as_ref().map(|_| self.sparse_io_gate.acquire());
1341 let sid = segment.meta().id;
1342 let (mut results, segment_seen) = crate::query::search_segment_shared_sync_planned(
1343 segment.as_ref(),
1344 query,
1345 fetch_limit,
1346 collect_positions,
1347 shared.clone(),
1348 lsp_plan,
1349 text_stats.clone(),
1350 )?;
1351 if fetch_limit > 0 && results.len() >= fetch_limit {
1352 shared.raise(results[fetch_limit - 1].score);
1353 }
1354 for result in &mut results {
1355 result.segment_id = sid;
1356 }
1357 Ok::<_, crate::Error>((results, segment_seen))
1358 };
1359
1360 if !lsp_plans.iter().any(Option::is_some) {
1361 let (merged, seen) = self.install_search_cpu(|| {
1362 (0..self.segments.len())
1363 .into_par_iter()
1364 .map(|segment| run_segment(&segment))
1365 .try_reduce(
1366 || (Vec::new(), 0u32),
1367 |(left, left_seen), (right, right_seen)| {
1368 Ok((
1369 merge_two_ranked(left, right, fetch_limit),
1370 left_seen.saturating_add(right_seen),
1371 ))
1372 },
1373 )
1374 })?;
1375 return Ok((merged, seen, shared.truncated()));
1376 }
1377
1378 let order = self.ordered_lsp_segments(&lsp_plans);
1379
1380 let mut merged = Vec::new();
1381 let mut merge_scratch = Vec::new();
1382 let mut total_seen = 0u32;
1383 if let Some((&pilot, rest)) = order.split_first() {
1384 let (pilot_results, pilot_seen) = self.search_pool.install(|| run_segment(&pilot))?;
1385 merge_ranked_reuse(&mut merged, pilot_results, fetch_limit, &mut merge_scratch);
1386 total_seen = total_seen.saturating_add(pilot_seen);
1387
1388 for wave in rest.chunks(self.bmp_wave_width()) {
1389 let batches = self
1390 .search_pool
1391 .install(|| wave.par_iter().map(run_segment).collect::<Result<Vec<_>>>())?;
1392 for (results, seen) in batches {
1393 merge_ranked_reuse(&mut merged, results, fetch_limit, &mut merge_scratch);
1394 total_seen = total_seen.saturating_add(seen);
1395 }
1396 }
1397 }
1398 Ok((merged, total_seen, shared.truncated()))
1399 }
1400
1401 #[cfg(feature = "sync")]
1405 pub fn search_with_offset_and_count_sync(
1406 &self,
1407 query: &dyn crate::query::Query,
1408 limit: usize,
1409 offset: usize,
1410 ) -> Result<(Vec<crate::query::SearchResult>, u32)> {
1411 let fetch_limit = checked_search_window(limit, offset)?;
1412 let text_stats = self.query_text_stats(query, None);
1413 let (merged, total_seen, _) =
1414 self.search_segments_sync(query, fetch_limit, false, None, text_stats)?;
1415
1416 let results = apply_result_offset(merged, fetch_limit, offset);
1417 Ok((results, total_seen))
1418 }
1419
1420 pub async fn search_fused(
1441 &self,
1442 queries: &[(&dyn crate::query::Query, f32)],
1443 fetch_limit: usize,
1444 limit: usize,
1445 method: crate::query::FusionMethod,
1446 combiner: crate::query::MultiValueCombiner,
1447 ) -> Result<Vec<crate::query::SearchResult>> {
1448 let (results, _) = self
1449 .search_fused_with_count(queries, fetch_limit, limit, method, combiner)
1450 .await?;
1451 Ok(results)
1452 }
1453
1454 pub async fn search_fused_with_count(
1458 &self,
1459 queries: &[(&dyn crate::query::Query, f32)],
1460 fetch_limit: usize,
1461 limit: usize,
1462 method: crate::query::FusionMethod,
1463 combiner: crate::query::MultiValueCombiner,
1464 ) -> Result<(Vec<crate::query::SearchResult>, u32)> {
1465 if queries.is_empty() {
1466 return Err(crate::Error::Query(
1467 "fusion requires at least one sub-query".to_string(),
1468 ));
1469 }
1470 if queries.len() > crate::query::MAX_FUSION_SUB_QUERIES {
1471 return Err(crate::Error::Query(format!(
1472 "fusion supports at most {} sub-queries, got {}",
1473 crate::query::MAX_FUSION_SUB_QUERIES,
1474 queries.len()
1475 )));
1476 }
1477 if fetch_limit == 0 {
1478 return Err(crate::Error::Query(
1479 "fusion fetch_limit must be greater than zero".to_string(),
1480 ));
1481 }
1482 let candidate_slots = fetch_limit
1483 .checked_mul(queries.len())
1484 .ok_or_else(|| crate::Error::Query("fusion candidate budget overflow".to_string()))?;
1485 if candidate_slots > crate::query::MAX_FUSION_CANDIDATE_SLOTS {
1486 return Err(crate::Error::Query(format!(
1487 "fusion candidate budget must not exceed {}, got {candidate_slots}",
1488 crate::query::MAX_FUSION_CANDIDATE_SLOTS
1489 )));
1490 }
1491 for (index, &(_, weight)) in queries.iter().enumerate() {
1492 if !weight.is_finite() || weight < 0.0 {
1493 return Err(crate::Error::Query(format!(
1494 "fusion query weight at index {index} must be finite and non-negative, \
1495 got {weight}"
1496 )));
1497 }
1498 }
1499 if let crate::query::FusionMethod::Rrf { k } = method
1500 && (!k.is_finite() || k < 0.0)
1501 {
1502 return Err(crate::Error::Query(format!(
1503 "fusion RRF k must be finite and non-negative, got {k}"
1504 )));
1505 }
1506 combiner.validate().map_err(crate::Error::Query)?;
1507
1508 #[cfg(feature = "sync")]
1521 if !self.segments.is_empty()
1522 && tokio::runtime::Handle::current().runtime_flavor()
1523 == tokio::runtime::RuntimeFlavor::MultiThread
1524 {
1525 use rayon::prelude::*;
1526 let lists: Vec<(Vec<crate::query::SearchResult>, f32, u32)> =
1527 tokio::task::block_in_place(|| {
1528 queries
1529 .par_iter()
1530 .map(|&(query, weight)| {
1531 let (results, seen) =
1532 self.search_internal_sync(query, fetch_limit, 0, true)?;
1533 Ok((results, weight, seen))
1534 })
1535 .collect::<Result<Vec<_>>>()
1536 })?;
1537 let mut total_seen = 0u32;
1538 let ranked_lists = lists
1539 .into_iter()
1540 .map(|(results, weight, seen)| {
1541 total_seen = total_seen.saturating_add(seen);
1542 (results, weight)
1543 })
1544 .collect();
1545 let fused =
1546 crate::query::try_fuse_ranked_lists_chunked(ranked_lists, method, combiner, limit)
1547 .map_err(crate::Error::Query)?;
1548 return Ok((fused, total_seen));
1549 }
1550
1551 let lists: Vec<(Vec<crate::query::SearchResult>, f32, u32)> =
1554 futures::future::try_join_all(queries.iter().map(|&(query, weight)| async move {
1555 let (results, seen) = self.search_with_positions(query, fetch_limit).await?;
1556 Ok::<_, crate::Error>((results, weight, seen))
1557 }))
1558 .await?;
1559 let mut total_seen = 0u32;
1560 let ranked_lists = lists
1561 .into_iter()
1562 .map(|(results, weight, seen)| {
1563 total_seen = total_seen.saturating_add(seen);
1564 (results, weight)
1565 })
1566 .collect();
1567 let fused =
1568 crate::query::try_fuse_ranked_lists_chunked(ranked_lists, method, combiner, limit)
1569 .map_err(crate::Error::Query)?;
1570 Ok((fused, total_seen))
1571 }
1572
1573 pub fn fuse_candidate_lists(
1575 &self,
1576 lists: &[(&[crate::query::SearchResult], f32)],
1577 method: crate::query::FusionMethod,
1578 combiner: crate::query::MultiValueCombiner,
1579 limit: usize,
1580 ) -> Result<Vec<crate::query::SearchResult>> {
1581 self.install_search_cpu(|| {
1582 crate::query::try_fuse_ranked_lists_chunked_borrowed(lists, method, combiner, limit)
1583 })
1584 .map_err(crate::Error::Query)
1585 }
1586
1587 pub fn rrf_scores_for_hits(
1589 &self,
1590 lists: &[crate::query::RrfRankedList<'_>],
1591 selected: &[(u128, u32)],
1592 k: f32,
1593 combiner: crate::query::MultiValueCombiner,
1594 ) -> Result<Vec<crate::query::RrfScore>> {
1595 self.install_search_cpu(|| crate::query::rrf_scores_for_hits(lists, selected, k, combiner))
1596 .map_err(crate::Error::Query)
1597 }
1598
1599 pub async fn search_candidate_union(
1602 &self,
1603 queries: &[Arc<dyn crate::query::Query>],
1604 depth: usize,
1605 stats: Arc<crate::query::GlobalStats>,
1606 ) -> Result<(Vec<crate::query::SearchResult>, u32)> {
1607 let lists = self
1608 .search_candidate_lists(queries, depth, Some(stats))
1609 .await?;
1610 let total_seen = lists
1611 .iter()
1612 .fold(0u32, |sum, (_, seen)| sum.saturating_add(*seen));
1613 let union = self.merge_candidate_lists(lists.into_iter().map(|(list, _)| list))?;
1614 Ok((union, total_seen))
1615 }
1616
1617 pub async fn search_candidate_lists(
1620 &self,
1621 queries: &[Arc<dyn crate::query::Query>],
1622 depth: usize,
1623 stats: Option<Arc<crate::query::GlobalStats>>,
1624 ) -> Result<Vec<(Vec<crate::query::SearchResult>, u32)>> {
1625 if queries.is_empty()
1626 || queries.len() > crate::query::MAX_FUSION_SUB_QUERIES
1627 || depth == 0
1628 || depth
1629 .checked_mul(queries.len())
1630 .is_none_or(|slots| slots > crate::query::MAX_FUSION_CANDIDATE_SLOTS)
1631 {
1632 return Err(crate::Error::Query(
1633 "invalid candidate union branch/depth budget".into(),
1634 ));
1635 }
1636 #[cfg(feature = "sync")]
1637 let lists = if !self.segments.is_empty()
1638 && tokio::runtime::Handle::current().runtime_flavor()
1639 == tokio::runtime::RuntimeFlavor::MultiThread
1640 {
1641 use rayon::prelude::*;
1642 tokio::task::block_in_place(|| {
1643 self.install_search_cpu(|| {
1644 queries
1645 .par_iter()
1646 .map(|query| {
1647 let (results, seen, _) = self.search_internal_sync_budgeted(
1648 query.as_ref(),
1649 depth,
1650 0,
1651 true,
1652 None,
1653 self.query_text_stats(query.as_ref(), stats.clone()),
1654 )?;
1655 Ok((results, seen))
1656 })
1657 .collect::<Result<Vec<_>>>()
1658 })
1659 })?
1660 } else {
1661 self.search_candidate_branches_async(queries, depth, stats)
1662 .await?
1663 };
1664 #[cfg(not(feature = "sync"))]
1665 let lists = self
1666 .search_candidate_branches_async(queries, depth, stats)
1667 .await?;
1668 Ok(lists)
1669 }
1670
1671 pub fn merge_candidate_lists<
1672 T: Into<crate::query::SearchResult> + std::borrow::Borrow<crate::query::SearchResult>,
1673 >(
1674 &self,
1675 lists: impl IntoIterator<Item = impl IntoIterator<Item = T>>,
1676 ) -> Result<Vec<crate::query::SearchResult>> {
1677 let mut union = std::collections::BTreeMap::new();
1678 let mut positions = 0usize;
1679 let mut slots = 0usize;
1680 for list in lists {
1681 for result in list {
1682 slots = slots.saturating_add(1);
1683 if slots > crate::query::MAX_FUSION_CANDIDATE_SLOTS {
1684 return Err(crate::Error::Query(
1685 "candidate union document budget exceeded".into(),
1686 ));
1687 }
1688 let borrowed = result.borrow();
1689 positions = positions.saturating_add(
1690 borrowed
1691 .positions
1692 .iter()
1693 .map(|(_, p)| p.len())
1694 .sum::<usize>(),
1695 );
1696 if positions > crate::query::MAX_FUSION_CHUNK_SLOTS {
1697 return Err(crate::Error::Query(
1698 "candidate union ordinal budget exceeded".into(),
1699 ));
1700 }
1701 let mut result: crate::query::SearchResult = result.into();
1704 result.score = 0.0;
1705 match union.entry((result.segment_id, result.doc_id)) {
1706 std::collections::btree_map::Entry::Vacant(entry) => {
1707 entry.insert(result);
1708 }
1709 std::collections::btree_map::Entry::Occupied(mut entry) => {
1710 entry.get_mut().positions.extend(result.positions);
1711 }
1712 }
1713 }
1714 }
1715 Ok(union.into_values().collect())
1716 }
1717
1718 async fn search_candidate_branches_async(
1719 &self,
1720 queries: &[Arc<dyn crate::query::Query>],
1721 depth: usize,
1722 stats: Option<Arc<crate::query::GlobalStats>>,
1723 ) -> Result<Vec<(Vec<crate::query::SearchResult>, u32)>> {
1724 futures::future::try_join_all(queries.iter().map(|query| {
1725 let stats = stats.clone();
1726 async move {
1727 let (results, seen, _) = self
1728 .search_with_positions_budgeted_stats(query.as_ref(), depth, None, stats)
1729 .await?;
1730 Ok((results, seen))
1731 }
1732 }))
1733 .await
1734 }
1735
1736 pub async fn search_and_rerank(
1741 &self,
1742 query: &dyn crate::query::Query,
1743 l1_limit: usize,
1744 final_limit: usize,
1745 config: &crate::query::RerankerConfig,
1746 ) -> Result<(Vec<crate::query::SearchResult>, u32)> {
1747 let (candidates, total_seen) = self.search_with_count(query, l1_limit).await?;
1748 let reranked = crate::query::rerank(self, &candidates, config, final_limit).await?;
1749 Ok((reranked, total_seen))
1750 }
1751
1752 pub async fn query(
1754 &self,
1755 query_str: &str,
1756 limit: usize,
1757 ) -> Result<crate::query::SearchResponse> {
1758 self.query_offset(query_str, limit, 0).await
1759 }
1760
1761 pub async fn query_offset(
1763 &self,
1764 query_str: &str,
1765 limit: usize,
1766 offset: usize,
1767 ) -> Result<crate::query::SearchResponse> {
1768 let parser = self.query_parser();
1769 let query = parser
1770 .parse(query_str)
1771 .map_err(crate::error::Error::Query)?;
1772
1773 let (results, _total_seen) = self
1774 .search_internal(query.as_ref(), limit, offset, false)
1775 .await?;
1776
1777 let total_hits = results.len() as u32;
1778 let hits: Vec<crate::query::SearchHit> = results
1779 .into_iter()
1780 .map(|result| crate::query::SearchHit {
1781 address: crate::query::DocAddress::new(result.segment_id, result.doc_id),
1782 score: result.score,
1783 matched_fields: result.extract_ordinals(),
1784 })
1785 .collect();
1786
1787 Ok(crate::query::SearchResponse { hits, total_hits })
1788 }
1789
1790 pub fn query_parser(&self) -> crate::dsl::QueryLanguageParser {
1792 let query_routers = self.schema.query_routers();
1793 if !query_routers.is_empty()
1794 && let Ok(router) = crate::dsl::QueryFieldRouter::from_rules(query_routers)
1795 {
1796 return crate::dsl::QueryLanguageParser::with_router(
1797 Arc::clone(&self.schema),
1798 self.default_fields.clone(),
1799 Arc::clone(&self.tokenizers),
1800 router,
1801 );
1802 }
1803
1804 crate::dsl::QueryLanguageParser::new(
1805 Arc::clone(&self.schema),
1806 self.default_fields.clone(),
1807 Arc::clone(&self.tokenizers),
1808 )
1809 }
1810
1811 pub async fn get_document(
1813 &self,
1814 address: &crate::query::DocAddress,
1815 ) -> Result<Option<crate::dsl::Document>> {
1816 self.get_document_with_fields(address, None).await
1817 }
1818
1819 pub async fn get_document_with_fields(
1825 &self,
1826 address: &crate::query::DocAddress,
1827 fields: Option<&rustc_hash::FxHashSet<u32>>,
1828 ) -> Result<Option<crate::dsl::Document>> {
1829 let segment_id = address.segment_id_u128().ok_or_else(|| {
1830 crate::error::Error::Query(format!("Invalid segment ID: {}", address.segment_id()))
1831 })?;
1832
1833 if let Some(&idx) = self.segment_map.get(&segment_id) {
1834 return self.segments[idx]
1835 .doc_with_fields(address.doc_id, fields)
1836 .await;
1837 }
1838
1839 Ok(None)
1840 }
1841}
1842
1843struct HierarchicalLspSelection {
1844 selected: Vec<Vec<(u32, f32)>>,
1845 expanded_groups: usize,
1846 evaluated_superblocks: usize,
1847}
1848
1849fn select_global_lsp_hierarchical(
1856 coarse_bounds: &[Option<Vec<f32>>],
1857 gamma: usize,
1858 mut expand: impl FnMut(usize, u32, &mut Vec<(u32, f32)>) -> Result<()>,
1859) -> Result<HierarchicalLspSelection> {
1860 if gamma == 0 {
1861 return Err(crate::Error::Internal(
1862 "hierarchical LSP selection requires a positive gamma".into(),
1863 ));
1864 }
1865 let mut frontier = std::collections::BinaryHeap::<(u32, usize, u32)>::new();
1866 for (segment, bounds) in coarse_bounds.iter().enumerate() {
1867 let Some(bounds) = bounds else {
1868 continue;
1869 };
1870 for (coarse_group, &bound) in bounds.iter().enumerate() {
1871 debug_assert!(bound.is_finite() && bound >= 0.0);
1872 if bound > 0.0 {
1873 frontier.push((bound.to_bits(), segment, coarse_group as u32));
1874 }
1875 }
1876 }
1877 let mut top =
1878 std::collections::BinaryHeap::<std::cmp::Reverse<(u32, usize, u32)>>::with_capacity(
1879 gamma.min(65_536),
1880 );
1881 let mut expanded =
1882 Vec::with_capacity(crate::segment::reader::bmp::BMP_COARSE_SUPERBLOCKS as usize);
1883 let mut expanded_groups = 0usize;
1884 let mut evaluated_superblocks = 0usize;
1885 while let Some((coarse_bound, segment, coarse_group)) = frontier.pop() {
1886 if top.len() == gamma && top.peek().is_some_and(|minimum| coarse_bound < minimum.0.0) {
1887 break;
1888 }
1889 expand(segment, coarse_group, &mut expanded)?;
1890 expanded_groups += 1;
1891 evaluated_superblocks = evaluated_superblocks.saturating_add(expanded.len());
1892 for &(superblock, bound) in &expanded {
1893 debug_assert!(bound.is_finite() && bound >= 0.0);
1894 if bound <= 0.0 {
1895 continue;
1896 }
1897 let candidate = (bound.to_bits(), segment, superblock);
1898 if top.len() < gamma {
1899 top.push(std::cmp::Reverse(candidate));
1900 } else if top.peek().is_some_and(|minimum| candidate > minimum.0) {
1901 top.pop();
1902 top.push(std::cmp::Reverse(candidate));
1903 }
1904 }
1905 }
1906
1907 let mut selected = vec![Vec::<(u32, f32)>::new(); coarse_bounds.len()];
1908 for std::cmp::Reverse((bound, segment, superblock)) in top {
1909 selected[segment].push((superblock, f32::from_bits(bound)));
1910 }
1911 for segment in &mut selected {
1912 segment.sort_unstable_by(|&(left_sb, left_bound), &(right_sb, right_bound)| {
1913 right_bound
1914 .total_cmp(&left_bound)
1915 .then_with(|| left_sb.cmp(&right_sb))
1916 });
1917 }
1918 Ok(HierarchicalLspSelection {
1919 selected,
1920 expanded_groups,
1921 evaluated_superblocks,
1922 })
1923}
1924
1925#[cfg(test)]
1929fn select_global_lsp_superblocks(bounds: &[Option<Vec<f32>>], gamma: usize) -> Vec<Vec<u32>> {
1930 let mut top =
1931 std::collections::BinaryHeap::<std::cmp::Reverse<(u32, usize, u32)>>::with_capacity(
1932 gamma.min(65_536),
1933 );
1934 for (segment, segment_bounds) in bounds.iter().enumerate() {
1935 let Some(segment_bounds) = segment_bounds else {
1936 continue;
1937 };
1938 for (superblock, &bound) in segment_bounds.iter().enumerate() {
1939 if bound <= 0.0 {
1940 continue;
1941 }
1942 let candidate = (bound.to_bits(), segment, superblock as u32);
1943 if top.len() < gamma {
1944 top.push(std::cmp::Reverse(candidate));
1945 } else if top.peek().is_some_and(|minimum| candidate > minimum.0) {
1946 top.pop();
1947 top.push(std::cmp::Reverse(candidate));
1948 }
1949 }
1950 }
1951
1952 let mut selected = vec![Vec::<u32>::new(); bounds.len()];
1953 for std::cmp::Reverse((_, segment, superblock)) in top {
1954 selected[segment].push(superblock);
1955 }
1956 selected
1957}
1958
1959fn merge_ranked_reuse(
1963 merged: &mut Vec<crate::query::SearchResult>,
1964 batch: Vec<crate::query::SearchResult>,
1965 limit: usize,
1966 scratch: &mut Vec<crate::query::SearchResult>,
1967) {
1968 scratch.clear();
1969 let output_len = limit.min(merged.len().saturating_add(batch.len()));
1970 if scratch.capacity() < output_len {
1971 scratch.reserve_exact(output_len);
1972 }
1973 {
1974 let mut left = merged.drain(..).peekable();
1975 let mut right = batch.into_iter().peekable();
1976 while scratch.len() < output_len {
1977 let take_left = match (left.peek(), right.peek()) {
1978 (Some(left), Some(right)) => {
1979 !crate::query::compare_search_results_desc(left, right).is_gt()
1980 }
1981 (Some(_), None) => true,
1982 (None, Some(_)) => false,
1983 (None, None) => break,
1984 };
1985 if take_left {
1986 scratch.push(left.next().expect("peeked left result"));
1987 } else {
1988 scratch.push(right.next().expect("peeked right result"));
1989 }
1990 }
1991 }
1992 std::mem::swap(merged, scratch);
1993}
1994
1995#[cfg(any(feature = "sync", test))]
1999fn merge_two_ranked(
2000 left: Vec<crate::query::SearchResult>,
2001 right: Vec<crate::query::SearchResult>,
2002 limit: usize,
2003) -> Vec<crate::query::SearchResult> {
2004 let mut left = left.into_iter().peekable();
2005 let mut right = right.into_iter().peekable();
2006 let mut merged = Vec::with_capacity(limit.min(left.len().saturating_add(right.len())));
2007
2008 while merged.len() < limit {
2009 let take_left = match (left.peek(), right.peek()) {
2010 (Some(left), Some(right)) => {
2011 !crate::query::compare_search_results_desc(left, right).is_gt()
2012 }
2013 (Some(_), None) => true,
2014 (None, Some(_)) => false,
2015 (None, None) => break,
2016 };
2017 if take_left {
2018 merged.push(left.next().expect("peeked left result"));
2019 } else {
2020 merged.push(right.next().expect("peeked right result"));
2021 }
2022 }
2023 merged
2024}
2025
2026fn apply_result_offset(
2027 mut results: Vec<crate::query::SearchResult>,
2028 fetch_limit: usize,
2029 offset: usize,
2030) -> Vec<crate::query::SearchResult> {
2031 if offset == 0 {
2032 results.truncate(fetch_limit);
2033 return results;
2034 }
2035 results
2039 .into_iter()
2040 .skip(offset)
2041 .take(fetch_limit.saturating_sub(offset))
2042 .collect()
2043}
2044
2045fn checked_search_window(limit: usize, offset: usize) -> Result<usize> {
2046 offset
2047 .checked_add(limit)
2048 .ok_or_else(|| crate::Error::Query("search offset + limit overflow".into()))
2049}
2050
2051fn process_rss_bytes() -> u64 {
2053 #[cfg(target_os = "linux")]
2054 {
2055 if let Ok(status) = std::fs::read_to_string("/proc/self/status") {
2057 for line in status.lines() {
2058 if let Some(rest) = line.strip_prefix("VmRSS:") {
2059 let kib: u64 = rest
2060 .trim()
2061 .trim_end_matches("kB")
2062 .trim()
2063 .parse()
2064 .unwrap_or(0);
2065 return kib.saturating_mul(1024);
2066 }
2067 }
2068 }
2069 0
2070 }
2071 #[cfg(target_os = "macos")]
2072 {
2073 use std::mem;
2075 #[repr(C)]
2076 struct TaskBasicInfo {
2077 virtual_size: u64,
2078 resident_size: u64,
2079 resident_size_max: u64,
2080 user_time: [u32; 2],
2081 system_time: [u32; 2],
2082 policy: i32,
2083 suspend_count: i32,
2084 }
2085 unsafe extern "C" {
2086 fn mach_task_self() -> u32;
2087 fn task_info(task: u32, flavor: u32, info: *mut TaskBasicInfo, count: *mut u32) -> i32;
2088 }
2089 const MACH_TASK_BASIC_INFO: u32 = 20;
2090 let mut info: TaskBasicInfo = unsafe { mem::zeroed() };
2091 let mut count = (mem::size_of::<TaskBasicInfo>() / mem::size_of::<u32>()) as u32;
2092 let ret = unsafe {
2093 task_info(
2094 mach_task_self(),
2095 MACH_TASK_BASIC_INFO,
2096 &mut info,
2097 &mut count,
2098 )
2099 };
2100 if ret == 0 { info.resident_size } else { 0 }
2101 }
2102 #[cfg(not(any(target_os = "linux", target_os = "macos")))]
2103 {
2104 0
2105 }
2106}
2107
2108#[cfg(test)]
2109mod load_segments_tests {
2110 use super::*;
2111
2112 #[cfg(feature = "native")]
2113 #[derive(Clone, Default)]
2114 struct SlowSegmentDirectory {
2115 inner: crate::directories::RamDirectory,
2116 active: Arc<std::sync::atomic::AtomicUsize>,
2117 peak: Arc<std::sync::atomic::AtomicUsize>,
2118 }
2119
2120 #[cfg(feature = "native")]
2121 #[async_trait::async_trait]
2122 impl Directory for SlowSegmentDirectory {
2123 async fn exists(&self, path: &std::path::Path) -> std::io::Result<bool> {
2124 self.inner.exists(path).await
2125 }
2126 async fn file_size(&self, path: &std::path::Path) -> std::io::Result<u64> {
2127 self.inner.file_size(path).await
2128 }
2129 async fn open_read(
2130 &self,
2131 path: &std::path::Path,
2132 ) -> std::io::Result<crate::directories::FileHandle> {
2133 use std::sync::atomic::Ordering::SeqCst;
2134 if path
2135 .extension()
2136 .is_some_and(|extension| extension == "meta")
2137 {
2138 struct ActiveRead<'a>(&'a std::sync::atomic::AtomicUsize);
2139 impl Drop for ActiveRead<'_> {
2140 fn drop(&mut self) {
2141 self.0.fetch_sub(1, SeqCst);
2142 }
2143 }
2144 let active = self.active.fetch_add(1, SeqCst) + 1;
2145 let _guard = ActiveRead(&self.active);
2146 self.peak.fetch_max(active, SeqCst);
2147 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
2148 }
2149 self.inner.open_read(path).await
2150 }
2151 async fn read_range(
2152 &self,
2153 path: &std::path::Path,
2154 range: std::ops::Range<u64>,
2155 ) -> std::io::Result<crate::directories::OwnedBytes> {
2156 self.inner.read_range(path, range).await
2157 }
2158 async fn list_files(
2159 &self,
2160 prefix: &std::path::Path,
2161 ) -> std::io::Result<Vec<std::path::PathBuf>> {
2162 self.inner.list_files(prefix).await
2163 }
2164 async fn open_lazy(
2165 &self,
2166 path: &std::path::Path,
2167 ) -> std::io::Result<crate::directories::FileHandle> {
2168 self.inner.open_lazy(path).await
2169 }
2170 }
2171
2172 #[cfg(feature = "native")]
2173 #[tokio::test]
2174 async fn opening_many_segments_bounds_inflight_reads_and_cancellation_releases_them() {
2175 use std::sync::atomic::Ordering::SeqCst;
2176 let directory = Arc::new(SlowSegmentDirectory::default());
2177 let mut builder = crate::Schema::builder();
2178 let text = builder.add_text_field("text", true, false);
2179 let schema = Arc::new(builder.build());
2180 let mut writer = crate::IndexWriter::create(
2181 directory.inner.clone(),
2182 (*schema).clone(),
2183 crate::IndexConfig {
2184 merge_policy: Box::new(crate::NoMergePolicy),
2185 ..Default::default()
2186 },
2187 )
2188 .await
2189 .unwrap();
2190 for _ in 0..8 {
2191 let mut document = crate::Document::new();
2192 document.add_text(text, "present");
2193 writer.add_document(document).unwrap();
2194 writer.commit().await.unwrap();
2195 }
2196 let metadata = crate::index::IndexMetadata::load(directory.as_ref())
2197 .await
2198 .unwrap();
2199 let ids = metadata.segment_ids();
2200 let searcher = Searcher::open(directory.clone(), schema.clone(), &ids, 8)
2201 .await
2202 .unwrap();
2203 assert_eq!(searcher.segment_readers().len(), 8);
2204 assert!(
2205 directory.peak.load(SeqCst) <= 2,
2206 "opening every segment concurrently multiplies reload scratch"
2207 );
2208 assert_eq!(directory.active.load(SeqCst), 0);
2209 assert_eq!(
2210 searcher
2211 .segment_readers()
2212 .iter()
2213 .map(|r| SegmentId(r.meta().id).to_hex())
2214 .collect::<Vec<_>>(),
2215 ids
2216 );
2217 let result = tokio::time::timeout(
2218 std::time::Duration::from_millis(1),
2219 Searcher::open(directory.clone(), schema, &ids, 8),
2220 )
2221 .await;
2222 assert!(result.is_err());
2223 assert_eq!(
2224 directory.active.load(SeqCst),
2225 0,
2226 "cancelled opens must release in-flight work"
2227 );
2228 }
2229
2230 #[tokio::test]
2231 async fn searcher_open_fails_loud_on_corrupt_metadata_segment_id() {
2232 let directory = Arc::new(crate::directories::RamDirectory::new());
2233 let schema = Arc::new(crate::dsl::SchemaBuilder::default().build());
2234
2235 let result =
2236 Searcher::open(directory, schema, &["not-a-hex-segment-id".to_string()], 8).await;
2237
2238 match result {
2239 Ok(searcher) => panic!(
2240 "corrupt segment ID must fail loud instead of silently serving {} segments",
2241 searcher.segment_readers().len()
2242 ),
2243 Err(crate::error::Error::Corruption(message)) => {
2244 assert!(message.contains("not-a-hex-segment-id"), "{message}");
2245 }
2246 Err(other) => panic!("expected Corruption error for invalid segment ID, got: {other}"),
2247 }
2248 }
2249}
2250
2251#[cfg(test)]
2252mod search_window_tests {
2253 use super::{
2254 apply_result_offset, checked_search_window, merge_ranked_reuse, merge_two_ranked,
2255 select_global_lsp_hierarchical, select_global_lsp_superblocks,
2256 };
2257 use crate::query::SearchResult;
2258
2259 fn result(segment_id: u128, doc_id: u32, score: f32) -> SearchResult {
2260 SearchResult {
2261 doc_id,
2262 score,
2263 segment_id,
2264 positions: Vec::new(),
2265 }
2266 }
2267
2268 #[test]
2269 fn search_window_is_checked() {
2270 assert_eq!(checked_search_window(7, 5).unwrap(), 12);
2271 assert!(checked_search_window(1, usize::MAX).is_err());
2272 }
2273
2274 #[test]
2275 fn bounded_merge_preserves_canonical_order_and_ties() {
2276 let left = vec![result(2, 9, 10.0), result(2, 3, 7.0)];
2277 let right = vec![result(1, 8, 10.0), result(1, 2, 7.0)];
2278
2279 let expected = merge_two_ranked(left.clone(), right.clone(), 3);
2280 let mut merged = left;
2281 let mut scratch = Vec::new();
2282 merge_ranked_reuse(&mut merged, right, 3, &mut scratch);
2283 assert_eq!(merged, expected);
2284 assert!(scratch.is_empty());
2285 let keys: Vec<_> = merged
2286 .iter()
2287 .map(|result| (result.score, result.segment_id, result.doc_id))
2288 .collect();
2289 assert_eq!(keys, vec![(10.0, 1, 8), (10.0, 2, 9), (7.0, 1, 2)]);
2290 }
2291
2292 #[test]
2293 fn result_offset_returns_only_the_requested_window() {
2294 let results = (0..8)
2295 .map(|doc_id| result(1, doc_id, 8.0 - doc_id as f32))
2296 .collect();
2297
2298 let page = apply_result_offset(results, 5, 2);
2299 assert_eq!(
2300 page.iter().map(|result| result.doc_id).collect::<Vec<_>>(),
2301 vec![2, 3, 4]
2302 );
2303 }
2304
2305 #[test]
2306 fn lsp_gamma_is_global_not_per_segment() {
2307 let bounds = vec![
2308 Some(vec![9.0, 1.0, 8.0]),
2309 Some(vec![7.0, 6.0, 0.0]),
2310 None,
2311 Some(vec![5.0, 4.0]),
2312 ];
2313 let mut selected = select_global_lsp_superblocks(&bounds, 4);
2314 for segment in &mut selected {
2315 segment.sort_unstable();
2316 }
2317 assert_eq!(selected.iter().map(Vec::len).sum::<usize>(), 4);
2318 assert_eq!(selected[0], vec![0, 2]);
2319 assert_eq!(selected[1], vec![0, 1]);
2320 assert!(selected[2].is_empty());
2321 assert!(selected[3].is_empty());
2322 }
2323
2324 #[test]
2325 fn hierarchical_lsp_matches_full_e_scan_exactly() {
2326 const GROUP: usize = crate::segment::reader::bmp::BMP_COARSE_SUPERBLOCKS as usize;
2327 let full = vec![
2328 Some(
2329 (0..700)
2330 .map(|index| 1_000.0 - index as f32)
2331 .collect::<Vec<_>>(),
2332 ),
2333 Some(
2334 (0..530)
2335 .map(|index| 995.0 - index as f32 * 1.25)
2336 .collect::<Vec<_>>(),
2337 ),
2338 None,
2339 ];
2340 let coarse: Vec<Option<Vec<f32>>> = full
2341 .iter()
2342 .map(|segment| {
2343 segment.as_ref().map(|bounds| {
2344 bounds
2345 .chunks(GROUP)
2346 .map(|group| group.iter().copied().fold(0.0, f32::max))
2347 .collect()
2348 })
2349 })
2350 .collect();
2351 let gamma = 37;
2352 let hierarchy = select_global_lsp_hierarchical(&coarse, gamma, |segment, group, output| {
2353 output.clear();
2354 let Some(bounds) = &full[segment] else {
2355 return Ok(());
2356 };
2357 let start = group as usize * GROUP;
2358 output.extend(
2359 bounds[start..(start + GROUP).min(bounds.len())]
2360 .iter()
2361 .enumerate()
2362 .map(|(within, &bound)| ((start + within) as u32, bound)),
2363 );
2364 Ok(())
2365 })
2366 .unwrap();
2367 let mut expected = select_global_lsp_superblocks(&full, gamma);
2368 for segment in &mut expected {
2369 segment.sort_unstable();
2370 }
2371 let mut actual: Vec<Vec<u32>> = hierarchy
2372 .selected
2373 .iter()
2374 .map(|segment| {
2375 let mut ids: Vec<_> = segment.iter().map(|&(id, _)| id).collect();
2376 ids.sort_unstable();
2377 ids
2378 })
2379 .collect();
2380 for segment in &mut actual {
2381 segment.sort_unstable();
2382 }
2383 assert_eq!(actual, expected);
2384 assert_eq!(actual.iter().map(Vec::len).sum::<usize>(), gamma);
2385 assert!(
2386 hierarchy.evaluated_superblocks
2387 < full.iter().filter_map(Option::as_ref).map(Vec::len).sum()
2388 );
2389 }
2390}
2391
2392#[cfg(test)]
2393mod fusion_parallelism_tests {
2394 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2395 async fn nomination_lists_inherit_the_same_global_bm25_statistics_as_search() {
2396 use crate::query::{BooleanQuery, Query, TermQuery};
2397 use std::sync::Arc;
2398 let directory = crate::directories::RamDirectory::new();
2399 let mut builder = crate::Schema::builder();
2400 let text = builder.add_text_field("text", true, false);
2401 let config = crate::IndexConfig {
2402 merge_policy: Box::new(crate::NoMergePolicy),
2403 ..Default::default()
2404 };
2405 let mut writer =
2406 crate::IndexWriter::create(directory.clone(), builder.build(), config.clone())
2407 .await
2408 .unwrap();
2409 for segment in [
2410 ["needle needle", "needle common", "common common"],
2411 ["needle", "common", "common rare rare"],
2412 ] {
2413 for value in segment {
2414 let mut doc = crate::Document::new();
2415 doc.add_text(text, value);
2416 writer.add_document(doc).unwrap();
2417 }
2418 writer.commit().await.unwrap();
2419 }
2420 let index = crate::Index::open(directory, config).await.unwrap();
2421 let reader = index.reader().await.unwrap();
2422 let searcher = reader.searcher().await.unwrap();
2423 let query: Arc<dyn Query> = Arc::new(
2424 BooleanQuery::new()
2425 .should(TermQuery::text(text, "needle"))
2426 .should(TermQuery::text(text, "common")),
2427 );
2428 let expected = searcher
2429 .search_with_positions(query.as_ref(), 6)
2430 .await
2431 .unwrap();
2432 let actual = searcher
2433 .search_candidate_lists(&[query], 6, None)
2434 .await
2435 .unwrap();
2436 assert_eq!(actual[0], expected);
2437 }
2438
2439 #[test]
2440 fn fusion_runs_subqueries_concurrently() {
2441 let source = include_str!("searcher.rs");
2442 let body = source
2443 .split("pub async fn search_fused_with_count")
2444 .nth(1)
2445 .and_then(|tail| tail.split("/// Two-stage search").next())
2446 .expect("bounded fusion search implementation");
2447
2448 assert!(
2452 body.contains(".par_iter()"),
2453 "sync fusion path must run sub-queries in parallel over the rayon pool"
2454 );
2455 assert!(
2456 body.contains("try_join_all"),
2457 "async fusion fallback must run sub-queries concurrently"
2458 );
2459 }
2460}