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::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    /// Shared whole-pass concurrency limit. Tests and low-level callers may
121    /// omit it; SegmentManager always supplies the application-wide gate.
122    reorder_permits: Option<Arc<tokio::sync::Semaphore>>,
123}
124
125impl SegmentMerger {
126    pub fn new(schema: Arc<Schema>) -> Self {
127        Self {
128            schema,
129            reorder_bmp: false,
130            background_pool: None,
131            granularity: crate::segment::reorder::BpGranularity::Auto,
132            bp_budget: crate::segment::BpBudget::full(),
133            bp_memory_budget: crate::segment::reorder::DEFAULT_MEMORY_BUDGET,
134            reorder_permits: None,
135        }
136    }
137
138    /// Enable BP reordering of BMP fields during the merge (see `reorder_bmp`).
139    pub fn with_bmp_reorder(mut self, reorder: bool) -> Self {
140        self.reorder_bmp = reorder;
141        self
142    }
143
144    /// Run merge-time BP on this bounded pool instead of the global one.
145    pub fn with_background_pool(mut self, pool: Option<Arc<rayon::ThreadPool>>) -> Self {
146        self.background_pool = pool;
147        self
148    }
149
150    /// Set merge-time BP granularity (see `granularity`).
151    pub fn with_granularity(mut self, granularity: crate::segment::reorder::BpGranularity) -> Self {
152        self.granularity = granularity;
153        self
154    }
155
156    /// Bound merge-time BP wall clock (see `bp_budget`).
157    pub fn with_bp_budget(mut self, budget: crate::segment::BpBudget) -> Self {
158        self.bp_budget = budget;
159        self
160    }
161
162    /// Memory budget for the BP forward index (see `bp_memory_budget`).
163    pub fn with_bp_memory_budget(mut self, bytes: usize) -> Self {
164        self.bp_memory_budget = bytes;
165        self
166    }
167
168    /// Share the application-wide whole-segment reorder gate.
169    pub fn with_reorder_permits(mut self, permits: Arc<tokio::sync::Semaphore>) -> Self {
170        self.reorder_permits = Some(permits);
171        self
172    }
173
174    /// Merge segments into one, streaming postings/positions/store directly to files.
175    ///
176    /// If `trained` is provided, dense vectors use O(1) cluster merge when possible
177    /// (compatible IVF-PQ), otherwise rebuilds ANN from global artifacts.
178    /// Without trained structures, only flat vectors are merged.
179    ///
180    /// Uses streaming writers so postings, positions, and store data flow directly
181    /// to files instead of buffering everything in memory. Only the term dictionary
182    /// (compact key+TermInfo entries) is buffered.
183    pub async fn merge<D: Directory + DirectoryWriter>(
184        &self,
185        dir: &D,
186        segments: &[SegmentReader],
187        new_segment_id: SegmentId,
188        trained: Option<&TrainedVectorStructures>,
189    ) -> Result<(SegmentMeta, MergeStats)> {
190        // Reject an unrepresentable merge before creating any output files.
191        // The previous late check left a complete orphan output behind after
192        // doing all expensive phases.
193        let total_docs: u32 = segments
194            .iter()
195            .try_fold(0u32, |acc, segment| acc.checked_add(segment.num_docs()))
196            .ok_or_else(|| {
197                crate::Error::Internal(format!(
198                    "Total document count exceeds u32::MAX ({})",
199                    u32::MAX
200                ))
201            })?;
202
203        let mut stats = MergeStats::default();
204        let files = SegmentFiles::new(new_segment_id.0);
205
206        // === Two-stage merge to bound page cache pressure ===
207        //
208        // Stage 1: postings + store + fast_fields (concurrent)
209        //   Touches .term_dict, .postings, .positions, .store, .fast files.
210        //
211        // Stage 2: sparse + dense vectors. Block-copy sparse work runs with
212        // dense vectors; BP sparse work runs first to bound peak memory.
213        //   Touches .sparse, .vectors files.
214        //
215        // Running all phases concurrently caused OOM on large merges because
216        // mmap'd source files from all 16+ segments compete for page cache
217        // simultaneously (200+ GB of mmap'd data for BMP grids alone).
218        // Two stages halve the concurrent working set.
219        let merge_start = std::time::Instant::now();
220
221        // ── Stage 1: text + store + fast fields ─────────────────────────
222        let postings_fut = async {
223            let mut postings_writer =
224                OffsetWriter::new(dir.streaming_writer_cold(&files.postings).await?);
225            let mut positions_writer =
226                OffsetWriter::new(dir.streaming_writer_cold(&files.positions).await?);
227            let mut term_dict_writer =
228                OffsetWriter::new(dir.streaming_writer_cold(&files.term_dict).await?);
229
230            let terms_processed = self
231                .merge_postings(
232                    segments,
233                    &mut term_dict_writer,
234                    &mut postings_writer,
235                    &mut positions_writer,
236                )
237                .await?;
238
239            let postings_bytes = postings_writer.offset() as usize;
240            let term_dict_bytes = term_dict_writer.offset() as usize;
241            let positions_bytes = positions_writer.offset();
242
243            postings_writer.finish()?;
244            term_dict_writer.finish()?;
245            if positions_bytes > 0 {
246                positions_writer.finish()?;
247            } else {
248                drop(positions_writer);
249                let _ = dir.delete(&files.positions).await;
250            }
251            log::info!(
252                "[merge] postings done: {} terms, term_dict={}, postings={}, positions={}",
253                terms_processed,
254                format_bytes(term_dict_bytes),
255                format_bytes(postings_bytes),
256                format_bytes(positions_bytes as usize),
257            );
258            Ok::<(usize, usize, usize), crate::Error>((
259                terms_processed,
260                term_dict_bytes,
261                postings_bytes,
262            ))
263        };
264
265        let store_fut = async {
266            let mut store_writer =
267                OffsetWriter::new(dir.streaming_writer_cold(&files.store).await?);
268            let store_num_docs = self.merge_store(segments, &mut store_writer).await?;
269            let bytes = store_writer.offset() as usize;
270            store_writer.finish()?;
271            Ok::<(usize, u32), crate::Error>((bytes, store_num_docs))
272        };
273
274        let fast_fut = async { self.merge_fast_fields(dir, segments, &files).await };
275
276        let (postings_result, store_result, fast_bytes) =
277            tokio::try_join!(postings_fut, store_fut, fast_fut)?;
278
279        log::info!(
280            "[merge] stage 1 done in {:.1}s (postings + store + fast)",
281            merge_start.elapsed().as_secs_f64()
282        );
283
284        // ── Stage 2: sparse + dense vectors ─────────────────────────────
285        // Page cache from stage 1 files can now be evicted by the kernel
286        // as stage 2 accesses different mmap regions (.sparse, .vectors).
287        let sparse_fut = async { self.merge_sparse_vectors(dir, segments, &files).await };
288
289        let dense_fut = async {
290            self.merge_dense_vectors(dir, segments, &files, trained, AnnWriteMode::Copy)
291                .await
292        };
293
294        // Merge-time BP constructs a potentially budget-sized forward index.
295        // Do not overlap that allocation and its heavy source-file scan with
296        // an ANN rebuild. Block-copy sparse merges remain concurrent with ANN.
297        let ((sparse_bytes, bp_converged), vectors_bytes) = if self.reorder_bmp {
298            let sparse = sparse_fut.await?;
299            let dense = dense_fut.await?;
300            (sparse, dense)
301        } else {
302            tokio::try_join!(sparse_fut, dense_fut)?
303        };
304        let (store_bytes, store_num_docs) = store_result;
305        stats.terms_processed = postings_result.0;
306        stats.term_dict_bytes = postings_result.1;
307        stats.postings_bytes = postings_result.2;
308        stats.store_bytes = store_bytes;
309        stats.vectors_bytes = vectors_bytes;
310        stats.sparse_bytes = sparse_bytes;
311        stats.bp_converged = bp_converged;
312        stats.fast_bytes = fast_bytes;
313        log::info!(
314            "[merge] all phases done in {:.1}s: {}",
315            merge_start.elapsed().as_secs_f64(),
316            stats
317        );
318
319        // === Mandatory: merge field stats + write meta ===
320        let mut merged_field_stats: FxHashMap<u32, FieldStats> = FxHashMap::default();
321        for segment in segments {
322            for (&field_id, field_stats) in &segment.meta().field_stats {
323                let entry = merged_field_stats.entry(field_id).or_default();
324                entry.total_tokens = entry
325                    .total_tokens
326                    .checked_add(field_stats.total_tokens)
327                    .ok_or_else(|| {
328                        crate::Error::Corruption(format!(
329                            "field {} total-token count overflow while merging",
330                            field_id
331                        ))
332                    })?;
333                entry.doc_count = entry
334                    .doc_count
335                    .checked_add(field_stats.doc_count)
336                    .ok_or_else(|| {
337                        crate::Error::Corruption(format!(
338                            "field {} document count overflow while merging",
339                            field_id
340                        ))
341                    })?;
342            }
343        }
344
345        // Verify store doc count matches metadata — a mismatch here means
346        // some store blocks were lost (e.g., compression thread panic) or
347        // source segment metadata disagrees with its store.
348        if store_num_docs != total_docs {
349            log::error!(
350                "[merge] STORE/META MISMATCH: store has {} docs but metadata expects {}. \
351                 Per-segment: {:?}",
352                store_num_docs,
353                total_docs,
354                segments
355                    .iter()
356                    .map(|s| (
357                        format!("{:016x}", s.meta().id),
358                        s.num_docs(),
359                        s.store().num_docs()
360                    ))
361                    .collect::<Vec<_>>()
362            );
363            return Err(crate::Error::Io(std::io::Error::new(
364                std::io::ErrorKind::InvalidData,
365                format!(
366                    "Store/meta doc count mismatch: store={}, meta={}",
367                    store_num_docs, total_docs
368                ),
369            )));
370        }
371
372        let meta = SegmentMeta {
373            id: new_segment_id.0,
374            num_docs: total_docs,
375            field_stats: merged_field_stats,
376        };
377
378        // Durable: replace_segments deletes the fsynced source segments right
379        // after publishing this output, so a non-durable .meta could be the
380        // only copy of the merged documents across a power failure.
381        dir.write_durable(&files.meta, &meta.serialize()?).await?;
382
383        let label = if trained.is_some() {
384            "ANN merge"
385        } else {
386            "Merge"
387        };
388        log::info!("{} complete: {} docs, {}", label, total_docs, stats);
389
390        Ok((meta, stats))
391    }
392}
393
394/// Delete segment files from directory (all deletions run concurrently).
395pub async fn delete_segment<D: Directory + DirectoryWriter>(
396    dir: &D,
397    segment_id: SegmentId,
398) -> Result<()> {
399    let files = SegmentFiles::new(segment_id.0);
400    let paths = files.lifecycle_paths();
401    let results = futures::future::join_all(paths.iter().map(|path| dir.delete(path))).await;
402
403    // Missing files are expected for optional components and idempotent
404    // retries. Any other failure must be surfaced so cleanup is not falsely
405    // reported as successful; a later orphan sweep can retry remaining files.
406    for result in results {
407        if let Err(error) = result
408            && error.kind() != std::io::ErrorKind::NotFound
409        {
410            return Err(crate::Error::Io(error));
411        }
412    }
413    Ok(())
414}