Skip to main content

hermes_core/segment/merger/
mod.rs

1//! Segment merger for combining multiple segments
2
3mod dense;
4#[cfg(feature = "diagnostics")]
5mod diagnostics;
6mod fast_fields;
7mod postings;
8mod sparse;
9mod store;
10
11use std::sync::Arc;
12
13use rustc_hash::FxHashMap;
14
15use super::reader::SegmentReader;
16use super::types::{FieldStats, SegmentFiles, SegmentId, SegmentMeta};
17use super::{OffsetWriter, format_bytes};
18use crate::Result;
19use crate::directories::{Directory, DirectoryWriter};
20use crate::dsl::Schema;
21
22/// Compute per-segment doc ID offsets (each segment's docs start after the previous).
23///
24/// Returns an error if the total document count across segments exceeds `u32::MAX`.
25fn doc_offsets(segments: &[SegmentReader]) -> Result<Vec<u32>> {
26    let mut offsets = Vec::with_capacity(segments.len());
27    let mut acc = 0u32;
28    for seg in segments {
29        offsets.push(acc);
30        acc = acc.checked_add(seg.num_docs()).ok_or_else(|| {
31            crate::Error::Internal(format!(
32                "Total document count across segments exceeds u32::MAX ({})",
33                u32::MAX
34            ))
35        })?;
36    }
37    Ok(offsets)
38}
39
40/// Statistics for merge operations
41#[derive(Debug, Clone, Default)]
42pub struct MergeStats {
43    /// Number of terms processed
44    pub terms_processed: usize,
45    /// Term dictionary output size
46    pub term_dict_bytes: usize,
47    /// Postings output size
48    pub postings_bytes: usize,
49    /// Store output size
50    pub store_bytes: usize,
51    /// Vector index output size
52    pub vectors_bytes: usize,
53    /// Sparse vector index output size
54    pub sparse_bytes: usize,
55    /// Whether merge-time BP reorder ran to full depth on every BMP field
56    /// (false = a pass hit its wall-clock budget; the segment is valid and
57    /// better-ordered, and the background optimizer deepens it later).
58    /// True when no BP ran (block-copy merges have nothing to deepen... they
59    /// are simply not reordered and tracked by the `reordered` flag instead).
60    pub bp_converged: bool,
61    /// Fast-field output size
62    pub fast_bytes: usize,
63}
64
65impl std::fmt::Display for MergeStats {
66    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
67        write!(
68            f,
69            "terms={}, term_dict={}, postings={}, store={}, vectors={}, sparse={}, fast={}",
70            self.terms_processed,
71            format_bytes(self.term_dict_bytes),
72            format_bytes(self.postings_bytes),
73            format_bytes(self.store_bytes),
74            format_bytes(self.vectors_bytes),
75            format_bytes(self.sparse_bytes),
76            format_bytes(self.fast_bytes),
77        )
78    }
79}
80
81// TrainedVectorStructures is defined in super::types (available on all platforms)
82pub use super::types::TrainedVectorStructures;
83
84/// Run a CPU/IO-heavy synchronous section, telling tokio to migrate this
85/// worker's task queue first (multi-thread runtimes only — `block_in_place`
86/// panics on current_thread, where we just run inline).
87pub(crate) fn block_in_place_if_multithread<R>(f: impl FnOnce() -> R) -> R {
88    if tokio::runtime::Handle::try_current()
89        .map(|h| h.runtime_flavor() == tokio::runtime::RuntimeFlavor::MultiThread)
90        .unwrap_or(false)
91    {
92        tokio::task::block_in_place(f)
93    } else {
94        f()
95    }
96}
97
98/// Segment merger - merges multiple segments into one
99pub struct SegmentMerger {
100    schema: Arc<Schema>,
101    /// Run BP reordering on BMP sparse fields while writing the merged blob
102    /// (instead of byte-level block stacking). The output segment is then
103    /// already ordered, so the standalone reorder pass is unnecessary.
104    reorder_bmp: bool,
105    /// Bounded rayon pool for merge-time BP. `None` = global pool (tests);
106    /// the SegmentManager always passes its background pool so BP cannot
107    /// starve query scoring.
108    background_pool: Option<Arc<rayon::ThreadPool>>,
109    /// Granularity for merge-time BP. `Auto` by default; the SegmentManager
110    /// forces `Records` when any merge source is an unconverged partial
111    /// reorder.
112    granularity: crate::segment::reorder::BpGranularity,
113    /// Budget for merge-time BP. Default unbudgeted; the SegmentManager
114    /// passes the index's `merge_bp_time_budget` so huge merges stop holding
115    /// a merge slot for the full BP depth — a truncated pass is marked
116    /// `bp_converged = false` and the background optimizer deepens it.
117    bp_budget: crate::segment::BpBudget,
118    /// Memory budget for the BP forward index during merge-time reorder.
119    bp_memory_budget: usize,
120}
121
122impl SegmentMerger {
123    pub fn new(schema: Arc<Schema>) -> Self {
124        Self {
125            schema,
126            reorder_bmp: false,
127            background_pool: None,
128            granularity: crate::segment::reorder::BpGranularity::Auto,
129            bp_budget: crate::segment::BpBudget::full(),
130            bp_memory_budget: crate::segment::reorder::DEFAULT_MEMORY_BUDGET,
131        }
132    }
133
134    /// Enable BP reordering of BMP fields during the merge (see `reorder_bmp`).
135    pub fn with_bmp_reorder(mut self, reorder: bool) -> Self {
136        self.reorder_bmp = reorder;
137        self
138    }
139
140    /// Run merge-time BP on this bounded pool instead of the global one.
141    pub fn with_background_pool(mut self, pool: Option<Arc<rayon::ThreadPool>>) -> Self {
142        self.background_pool = pool;
143        self
144    }
145
146    /// Set merge-time BP granularity (see `granularity`).
147    pub fn with_granularity(mut self, granularity: crate::segment::reorder::BpGranularity) -> Self {
148        self.granularity = granularity;
149        self
150    }
151
152    /// Bound merge-time BP wall clock (see `bp_budget`).
153    pub fn with_bp_budget(mut self, budget: crate::segment::BpBudget) -> Self {
154        self.bp_budget = budget;
155        self
156    }
157
158    /// Memory budget for the BP forward index (see `bp_memory_budget`).
159    pub fn with_bp_memory_budget(mut self, bytes: usize) -> Self {
160        self.bp_memory_budget = bytes;
161        self
162    }
163
164    /// Merge segments into one, streaming postings/positions/store directly to files.
165    ///
166    /// If `trained` is provided, dense vectors use O(1) cluster merge when possible
167    /// (homogeneous IVF/ScaNN), otherwise rebuilds ANN from trained structures.
168    /// Without trained structures, only flat vectors are merged.
169    ///
170    /// Uses streaming writers so postings, positions, and store data flow directly
171    /// to files instead of buffering everything in memory. Only the term dictionary
172    /// (compact key+TermInfo entries) is buffered.
173    pub async fn merge<D: Directory + DirectoryWriter>(
174        &self,
175        dir: &D,
176        segments: &[SegmentReader],
177        new_segment_id: SegmentId,
178        trained: Option<&TrainedVectorStructures>,
179    ) -> Result<(SegmentMeta, MergeStats)> {
180        let mut stats = MergeStats::default();
181        let files = SegmentFiles::new(new_segment_id.0);
182
183        // === Two-stage merge to bound page cache pressure ===
184        //
185        // Stage 1: postings + store + fast_fields (concurrent)
186        //   Touches .term_dict, .postings, .positions, .store, .fast files.
187        //
188        // Stage 2: sparse + dense vectors (concurrent)
189        //   Touches .sparse, .vectors files.
190        //
191        // Running all phases concurrently caused OOM on large merges because
192        // mmap'd source files from all 16+ segments compete for page cache
193        // simultaneously (200+ GB of mmap'd data for BMP grids alone).
194        // Two stages halve the concurrent working set.
195        let merge_start = std::time::Instant::now();
196
197        // ── Stage 1: text + store + fast fields ─────────────────────────
198        let postings_fut = async {
199            let mut postings_writer =
200                OffsetWriter::new(dir.streaming_writer_cold(&files.postings).await?);
201            let mut positions_writer =
202                OffsetWriter::new(dir.streaming_writer_cold(&files.positions).await?);
203            let mut term_dict_writer =
204                OffsetWriter::new(dir.streaming_writer_cold(&files.term_dict).await?);
205
206            let terms_processed = self
207                .merge_postings(
208                    segments,
209                    &mut term_dict_writer,
210                    &mut postings_writer,
211                    &mut positions_writer,
212                )
213                .await?;
214
215            let postings_bytes = postings_writer.offset() as usize;
216            let term_dict_bytes = term_dict_writer.offset() as usize;
217            let positions_bytes = positions_writer.offset();
218
219            postings_writer.finish()?;
220            term_dict_writer.finish()?;
221            if positions_bytes > 0 {
222                positions_writer.finish()?;
223            } else {
224                drop(positions_writer);
225                let _ = dir.delete(&files.positions).await;
226            }
227            log::info!(
228                "[merge] postings done: {} terms, term_dict={}, postings={}, positions={}",
229                terms_processed,
230                format_bytes(term_dict_bytes),
231                format_bytes(postings_bytes),
232                format_bytes(positions_bytes as usize),
233            );
234            Ok::<(usize, usize, usize), crate::Error>((
235                terms_processed,
236                term_dict_bytes,
237                postings_bytes,
238            ))
239        };
240
241        let store_fut = async {
242            let mut store_writer =
243                OffsetWriter::new(dir.streaming_writer_cold(&files.store).await?);
244            let store_num_docs = self.merge_store(segments, &mut store_writer).await?;
245            let bytes = store_writer.offset() as usize;
246            store_writer.finish()?;
247            Ok::<(usize, u32), crate::Error>((bytes, store_num_docs))
248        };
249
250        let fast_fut = async { self.merge_fast_fields(dir, segments, &files).await };
251
252        let (postings_result, store_result, fast_bytes) =
253            tokio::try_join!(postings_fut, store_fut, fast_fut)?;
254
255        log::info!(
256            "[merge] stage 1 done in {:.1}s (postings + store + fast)",
257            merge_start.elapsed().as_secs_f64()
258        );
259
260        // ── Stage 2: sparse + dense vectors ─────────────────────────────
261        // Page cache from stage 1 files can now be evicted by the kernel
262        // as stage 2 accesses different mmap regions (.sparse, .vectors).
263        let sparse_fut = async { self.merge_sparse_vectors(dir, segments, &files).await };
264
265        let dense_fut = async {
266            self.merge_dense_vectors(dir, segments, &files, trained)
267                .await
268        };
269
270        let ((sparse_bytes, bp_converged), vectors_bytes) =
271            tokio::try_join!(sparse_fut, dense_fut)?;
272        let (store_bytes, store_num_docs) = store_result;
273        stats.terms_processed = postings_result.0;
274        stats.term_dict_bytes = postings_result.1;
275        stats.postings_bytes = postings_result.2;
276        stats.store_bytes = store_bytes;
277        stats.vectors_bytes = vectors_bytes;
278        stats.sparse_bytes = sparse_bytes;
279        stats.bp_converged = bp_converged;
280        stats.fast_bytes = fast_bytes;
281        log::info!(
282            "[merge] all phases done in {:.1}s: {}",
283            merge_start.elapsed().as_secs_f64(),
284            stats
285        );
286
287        // === Mandatory: merge field stats + write meta ===
288        let mut merged_field_stats: FxHashMap<u32, FieldStats> = FxHashMap::default();
289        for segment in segments {
290            for (&field_id, field_stats) in &segment.meta().field_stats {
291                let entry = merged_field_stats.entry(field_id).or_default();
292                entry.total_tokens += field_stats.total_tokens;
293                entry.doc_count += field_stats.doc_count;
294            }
295        }
296
297        let total_docs: u32 = segments
298            .iter()
299            .try_fold(0u32, |acc, s| acc.checked_add(s.num_docs()))
300            .ok_or_else(|| {
301                crate::Error::Internal(format!(
302                    "Total document count exceeds u32::MAX ({})",
303                    u32::MAX
304                ))
305            })?;
306
307        // Verify store doc count matches metadata — a mismatch here means
308        // some store blocks were lost (e.g., compression thread panic) or
309        // source segment metadata disagrees with its store.
310        if store_num_docs != total_docs {
311            log::error!(
312                "[merge] STORE/META MISMATCH: store has {} docs but metadata expects {}. \
313                 Per-segment: {:?}",
314                store_num_docs,
315                total_docs,
316                segments
317                    .iter()
318                    .map(|s| (
319                        format!("{:016x}", s.meta().id),
320                        s.num_docs(),
321                        s.store().num_docs()
322                    ))
323                    .collect::<Vec<_>>()
324            );
325            return Err(crate::Error::Io(std::io::Error::new(
326                std::io::ErrorKind::InvalidData,
327                format!(
328                    "Store/meta doc count mismatch: store={}, meta={}",
329                    store_num_docs, total_docs
330                ),
331            )));
332        }
333
334        let meta = SegmentMeta {
335            id: new_segment_id.0,
336            num_docs: total_docs,
337            field_stats: merged_field_stats,
338        };
339
340        dir.write(&files.meta, &meta.serialize()?).await?;
341
342        let label = if trained.is_some() {
343            "ANN merge"
344        } else {
345            "Merge"
346        };
347        log::info!("{} complete: {} docs, {}", label, total_docs, stats);
348
349        Ok((meta, stats))
350    }
351}
352
353/// Delete segment files from directory (all deletions run concurrently).
354pub async fn delete_segment<D: Directory + DirectoryWriter>(
355    dir: &D,
356    segment_id: SegmentId,
357) -> Result<()> {
358    let files = SegmentFiles::new(segment_id.0);
359    let paths = files.lifecycle_paths();
360    let results = futures::future::join_all(paths.iter().map(|path| dir.delete(path))).await;
361
362    // Missing files are expected for optional components and idempotent
363    // retries. Any other failure must be surfaced so cleanup is not falsely
364    // reported as successful; a later orphan sweep can retry remaining files.
365    for result in results {
366        if let Err(error) = result
367            && error.kind() != std::io::ErrorKind::NotFound
368        {
369            return Err(crate::Error::Io(error));
370        }
371    }
372    Ok(())
373}