Skip to main content

hermes_core/segment/merger/
mod.rs

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