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