Skip to main content

summa_core/index/
searcher.rs

1//! Searcher - read-only search over pre-built segments
2//!
3//! This module provides `Searcher` for read-only search access to indexes.
4//! It can be used standalone (for wasm/read-only) or via `IndexReader` (for native).
5
6use 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/// Immutable resources that must stay identical across `IndexReader` reloads.
19/// The search pool exists only when synchronous scoring is compiled in; native
20/// async-only builds retain the cache policy without spawning unused threads.
21#[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    /// Cache and CPU policy of an `Index` opened with `config`.
36    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    /// Validates every load-time limit once, before any segment or file is
47    /// touched: the dictionary block cap and
48    /// the CPU/I-O widths. Later per-segment constructors re-check only what
49    /// they own.
50    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
84/// Searcher - provides search over loaded segments
85///
86/// For wasm/read-only use, create via `Searcher::open()`.
87/// For native use with Index, this is created via `IndexReader`.
88pub struct Searcher<D: Directory + 'static> {
89    /// Segment snapshot holding refs - prevents deletion during native use
90    #[cfg(feature = "native")]
91    _snapshot: SegmentSnapshot,
92    /// PhantomData for the directory generic
93    _phantom: std::marker::PhantomData<D>,
94    /// Loaded segment readers
95    segments: Vec<Arc<SegmentReader>>,
96    /// Schema
97    schema: Arc<Schema>,
98    /// Default fields for search
99    default_fields: Vec<crate::Field>,
100    /// Tokenizers
101    tokenizers: Arc<crate::tokenizer::TokenizerRegistry>,
102    /// One immutable generation of all index-global ANN artifacts.
103    trained_vectors: Arc<TrainedVectorStructures>,
104    /// Lazy global statistics for cross-segment IDF computation
105    global_stats: Arc<LazyGlobalStats>,
106    /// O(1) segment lookup by segment_id
107    segment_map: FxHashMap<u128, usize>,
108    /// Total document count across all segments
109    total_docs: u32,
110    /// Bounded process-wide-by-width pool for the complete nested search tree.
111    #[cfg(feature = "sync")]
112    search_pool: Arc<rayon::ThreadPool>,
113    /// Shared random-I/O gate and per-query wave width for BMP.
114    #[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    /// Create a Searcher directly from segment IDs
122    ///
123    /// This is a simpler initialization path that doesn't require SegmentManager.
124    /// Use this for read-only access to pre-built indexes.
125    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    /// Create from a snapshot (for native IndexReader use)
150    #[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    /// Create from a snapshot, reusing existing segment readers for unchanged segments.
190    /// This avoids re-opening mmaps, fast fields, sparse indexes, etc. for segments
191    /// that weren't touched by merge.
192    #[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    /// Internal create method
233    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; // suppress unused warning on wasm
286        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    /// Common loading logic shared by create and from_snapshot
308    #[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    /// Load segment readers with bounded opening concurrency.
351    /// Reuses existing segment readers for unchanged segments when `existing_segments`
352    /// is non-empty — avoids re-opening mmaps, fast fields, sparse indexes, etc.
353    #[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        // Build lookup from existing segment readers for reuse
366        let existing_map: FxHashMap<u128, Arc<SegmentReader>> = existing_segments
367            .iter()
368            .map(|seg| (seg.meta().id, Arc::clone(seg)))
369            .collect();
370
371        // Parse segment IDs from metadata. A key that fails to parse means the
372        // metadata is corrupt; fail loud instead of silently serving results
373        // without that segment's documents (the merge path errors with
374        // Corruption on the same input — search must not disagree).
375        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        // A changed sidecar is a new visibility view, not a new payload.
386        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        // Segment opens retain their completed payloads, but validation/copy
420        // scratch must not multiply by every new segment during reload.
421        const MAX_CONCURRENT_SEGMENT_OPENS: usize = 2;
422        use futures::{StreamExt, TryStreamExt};
423        // Separate independent copied indexes' shared document-cache keys.
424        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        // Sort by original index to maintain deterministic ordering
479        loaded.sort_by_key(|(idx, _)| *idx);
480
481        let segments: Vec<Arc<SegmentReader>> = loaded.into_iter().map(|(_, seg)| seg).collect();
482
483        // Keep heap, file-backed address space, and pinned residency separate.
484        // Mapped bytes are not necessarily resident; process RSS is the
485        // authoritative whole-process residency measurement.
486        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        // Log process RSS if available (helps diagnose OOM)
526        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        // One ANN-health line per dense field across the whole index, so an
544        // operator reads leaf skew and extent fragmentation from N_fields
545        // lines instead of N_segments × N_fields open-time lines.
546        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    /// Build default fields from schema
578    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    /// Get the schema
593    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    /// Get segment readers
602    pub fn segment_readers(&self) -> &[Arc<SegmentReader>] {
603        &self.segments
604    }
605
606    /// Get default fields for search
607    pub fn default_fields(&self) -> &[crate::Field] {
608        &self.default_fields
609    }
610
611    /// Get tokenizer registry
612    pub fn tokenizers(&self) -> &crate::tokenizer::TokenizerRegistry {
613        &self.tokenizers
614    }
615
616    /// Get trained centroids
617    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    /// Get lazy global statistics for cross-segment IDF computation
628    pub fn global_stats(&self) -> &Arc<LazyGlobalStats> {
629        &self.global_stats
630    }
631
632    /// Build O(1) lookup tables from loaded segments
633    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    /// Get total document count across all segments
644    pub fn num_docs(&self) -> u32 {
645        self.total_docs
646    }
647
648    /// Get O(1) segment_id → index map (used by reranker)
649    pub fn segment_map(&self) -> &FxHashMap<u128, usize> {
650        &self.segment_map
651    }
652
653    /// Run a bounded piece of CPU work inside this index's shared search pool.
654    #[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    /// Async-only/WASM builds execute inline because Rayon is not available.
681    /// Keeping this overload free of `Send` bounds allows browser-backed file
682    /// handles, whose callbacks are deliberately thread-local, to be scored.
683    #[cfg(not(feature = "sync"))]
684    pub(crate) fn install_search_cpu<R>(&self, operation: impl FnOnce() -> R) -> R {
685        operation()
686    }
687
688    /// Keep a ready scoring pipeline on the shared CPU pool. Polls borrow the
689    /// original future; pending I/O releases the worker and cancellation never
690    /// leaves detached scoring work holding a reader or request permit.
691    #[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    /// Get number of segments
724    pub fn num_segments(&self) -> usize {
725        self.segments.len()
726    }
727
728    /// Get a document by (segment_id, local_doc_id)
729    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    /// Search across all segments and return aggregated results
737    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    /// Search across all segments and return (results, total_seen)
747    /// total_seen is the number of documents that were scored across all segments
748    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    /// Search with offset for pagination
757    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    /// Search with offset and return (results, total_seen)
770    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    /// Search with positions (ordinal tracking) and return (results, total_seen)
780    ///
781    /// Use this when you need per-ordinal scores for multi-valued fields.
782    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    /// `search_with_positions` under a wall-clock budget (anytime mode).
791    /// Returns `(results, seen, truncated)`; `truncated` is set when an
792    /// executor stopped scoring at the deadline, in which case the results
793    /// are the best found so far rather than the exact top-k.
794    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    /// `search_with_count` under a wall-clock budget; see
805    /// [`Self::search_with_positions_budgeted`].
806    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    /// [`Self::search_with_positions_budgeted`] with externally supplied
817    /// text statistics (a broker's cross-shard document frequencies); they
818    /// replace the searcher's own segment-aggregated statistics.
819    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    /// [`Self::search_with_count_budgeted`] with externally supplied text
831    /// statistics.
832    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    /// Text statistics a query scores with: the caller's override when
844    /// given, otherwise document frequencies, corpus sizes and average
845    /// lengths aggregated over every segment of this searcher (`None` for a
846    /// single segment, whose local statistics already are the whole).
847    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    /// Build the paper's single query-level top-γ superblock set, then project
869    /// it back onto segment-local plans.
870    ///
871    /// Treating every immutable segment as an independent LSP index would
872    /// multiply work by the segment count. The prepass retains one global γ
873    /// while preserving Summa's streaming segment architecture.
874    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            // A cap covering the whole index is exhaustive. Let each segment
958            // compute and traverse its local order once instead of building a
959            // query-global heap and retaining an all-superblock selection.
960            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    /// Internal search implementation
1110    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        // Use rayon + block_in_place for CPU-bound scoring (sync feature required).
1136        // Offloads the scoring loop from tokio workers so search doesn't starve
1137        // other async tasks. Works for any segment count (rayon degrades gracefully
1138        // to inline execution for a single segment).
1139        // Only works on multi-threaded tokio runtime (block_in_place panics on current_thread).
1140        #[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        // No segments, no sync feature, or current_thread runtime: use an
1156        // explicitly bounded async stream. Starting every segment at once can
1157        // retain `segments × top_k` results while the slowest I/O completes.
1158        const MAX_ASYNC_SEGMENT_SEARCHES: usize = 8;
1159        use futures::StreamExt;
1160        use futures::TryStreamExt;
1161        // Cross-segment top-k floor (see search_internal_sync). Concurrent
1162        // segments share it via an atomic; ordering is best-effort. The floor
1163        // carries the query window so executors with clamped heaps (segments
1164        // smaller than the window) can never publish an invalid floor.
1165        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            // Match the synchronous policy: score the highest-bound pilot
1219            // first so lower-bound async work starts with a useful theta.
1220            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    /// Multi-segment parallel search using rayon (CPU-bound scoring on thread pool).
1245    ///
1246    /// `block_in_place` tells tokio this worker is occupied so it can steal tasks.
1247    /// `rayon::par_iter` distributes segment scoring across the rayon thread pool.
1248    #[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    /// Sync body of the parallel search: rayon par_iter over segments.
1271    /// Callers must already be off the async reactor (block_in_place or a
1272    /// rayon/blocking thread) — safe to nest inside another par_iter
1273    /// (rayon work-stealing composes).
1274    #[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    /// Score all segments with one shared threshold.
1311    ///
1312    /// Ordinary queries retain full CPU parallelism. BMP uses a highest-bound
1313    /// pilot followed by bounded waves, and every active BMP segment also
1314    /// holds a process-wide random-I/O permit. This lets theta mature before
1315    /// lower-bound segments touch pageable D/payload pages and prevents
1316    /// concurrent queries from multiplying the wave width.
1317    #[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    /// Synchronous search across all segments using rayon for parallelism.
1402    ///
1403    /// This is the async-free boundary — no tokio involvement from here down.
1404    #[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    /// Hybrid search: run several queries independently and fuse their
1421    /// ranked lists (union) into a single top-`limit` result.
1422    ///
1423    /// Unlike [`Self::search_and_rerank`] — which can only re-score
1424    /// documents the first-stage query already found — fusion keeps
1425    /// documents found by *any* of the queries. Typical use is sparse
1426    /// (BM25/SPLADE) + dense vector hybrid retrieval with
1427    /// `FusionMethod::Rrf { k: 60.0 }`.
1428    ///
1429    /// Fusion happens at **chunk granularity**: per-ordinal scores are
1430    /// collected from each sub-query, fused per `(doc, ordinal)` key, then
1431    /// combined into a doc score with `combiner`
1432    /// (`MultiValueCombiner::Max` recommended — same-chunk corroboration
1433    /// across verticals compounds, scattered noise does not). Fused results
1434    /// carry per-chunk `positions`.
1435    ///
1436    /// Each query is paired with a weight scaling its contribution.
1437    /// `fetch_limit` is the per-query candidate depth. Request-facing adapters
1438    /// default to at most [`crate::query::MAX_CANDIDATE_OVERSUBSCRIPTION`] times
1439    /// the result window and never multiply an existing rerank pool again.
1440    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    /// Fusion variant that also returns the aggregate number of documents
1455    /// scored by all sub-queries. This lets request-facing callers use the
1456    /// parallel fusion path without rerunning sub-queries for observability.
1457    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        // A hybrid query's sub-queries are independent (different fields and
1509        // scorers), so run them concurrently instead of one after another. Each
1510        // still fans out across every segment; rayon work-stealing overlaps the
1511        // branches' I/O stalls and fills the cores one branch leaves idle, so a
1512        // hybrid costs ~max(branch) rather than the sum. Sequential execution
1513        // left that on the table — measured 10–40% over max on production.
1514        //
1515        // Safe to nest: sub-query parallelism sits over segment parallelism on
1516        // the same shared pool (bounded by its thread count, not multiplied),
1517        // and the sparse BMP random-I/O gate is acquired per branch — hybrid's
1518        // binary branch never takes it, so the two cannot serialize on it.
1519        // `par_iter` preserves input order for deterministic rank ties.
1520        #[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        // Async/current-thread fallback: run the sub-queries concurrently so
1552        // their awaits overlap, preserving input order for deterministic ties.
1553        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    /// Fuse retained nomination lists on this searcher's shared CPU pool.
1574    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    /// Attribute organic RRF votes without repeating retrieval or changing L1.
1588    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    /// Retrieve independent branches and retain the complete bounded document
1600    /// union. No rank fusion or cross-branch top-k runs before feature scoring.
1601    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    /// Independent nominated lists for coordinator-side global fusion.
1618    /// This uses the same bounded parallel retrieval as candidate backfill.
1619    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                // Nomination magnitudes are intentionally not a mixed-vertical
1702                // score. Export requests get features; linear supplies rank.
1703                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    /// Two-stage search: L1 retrieval + L2 dense vector reranking
1737    ///
1738    /// Runs the query to get `l1_limit` candidates, then reranks by exact
1739    /// dense vector distance and returns the top `final_limit` results.
1740    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    /// Parse query string and search (convenience method)
1753    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    /// Parse query string and search with offset (convenience method)
1762    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    /// Get query parser for this searcher
1791    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    /// Get a document by address (segment_id + local doc_id)
1812    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    /// Get a document by address, hydrating only the specified field IDs.
1820    ///
1821    /// If `fields` is `None`, all fields are hydrated (including dense vectors).
1822    /// If `fields` is `Some(set)`, only dense vector fields in the set are read
1823    /// from flat storage — skipping expensive mmap reads for unrequested vectors.
1824    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
1849/// Select the exact global top-gamma E cells through safe H upper bounds.
1850///
1851/// Best-first expansion stops only when the next H bound is strictly below
1852/// the current gamma-th E bound. Equality must still expand because the
1853/// deterministic `(score, segment, superblock)` tie-break can change
1854/// membership. Memory is O(number of H cells + gamma), never O(all E cells).
1855fn 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/// Select one query-level top-γ set without allocating one tuple per
1926/// superblock. Prepared BMP bounds are finite and non-negative, so their f32
1927/// bit patterns have the same order as their numeric values.
1928#[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
1959/// Merge one segment batch into the running top-k while alternating two
1960/// retained buffers. Search results are moved, not cloned, and segment fan-out
1961/// no longer allocates a fresh `limit`-sized vector for every merge.
1962fn 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/// Merge two canonically sorted batches while moving (not cloning) hits.
1996/// Synchronous parallel reductions use this eagerly, so retained cross-segment
1997/// results stay O(k) instead of O(number_of_segments × k).
1998#[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    // Pagination usually returns a small window from a much larger fetch.
2036    // Allocate that small result rather than retaining the fetch-sized backing
2037    // allocation through response serialization.
2038    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
2051/// Get current process RSS in bytes (best-effort, returns zero on failure).
2052fn process_rss_bytes() -> u64 {
2053    #[cfg(target_os = "linux")]
2054    {
2055        // Read from /proc/self/status — VmRSS line
2056        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 mach_task_self / task_info via raw syscall
2074        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        // Independent hybrid branches run concurrently — rayon over the shared
2449        // pool in the multi-thread sync path, joined futures in the async
2450        // fallback — so a hybrid query costs ~max(branch) instead of the sum.
2451        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}