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    /// Fast-field output size
56    pub fast_bytes: usize,
57}
58
59impl std::fmt::Display for MergeStats {
60    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
61        write!(
62            f,
63            "terms={}, term_dict={}, postings={}, store={}, vectors={}, sparse={}, fast={}",
64            self.terms_processed,
65            format_bytes(self.term_dict_bytes),
66            format_bytes(self.postings_bytes),
67            format_bytes(self.store_bytes),
68            format_bytes(self.vectors_bytes),
69            format_bytes(self.sparse_bytes),
70            format_bytes(self.fast_bytes),
71        )
72    }
73}
74
75// TrainedVectorStructures is defined in super::types (available on all platforms)
76pub use super::types::TrainedVectorStructures;
77
78/// Segment merger - merges multiple segments into one
79pub struct SegmentMerger {
80    schema: Arc<Schema>,
81    /// Run BP reordering on BMP sparse fields while writing the merged blob
82    /// (instead of byte-level block stacking). The output segment is then
83    /// already ordered, so the standalone reorder pass is unnecessary.
84    reorder_bmp: bool,
85}
86
87impl SegmentMerger {
88    pub fn new(schema: Arc<Schema>) -> Self {
89        Self {
90            schema,
91            reorder_bmp: false,
92        }
93    }
94
95    /// Enable BP reordering of BMP fields during the merge (see `reorder_bmp`).
96    pub fn with_bmp_reorder(mut self, reorder: bool) -> Self {
97        self.reorder_bmp = reorder;
98        self
99    }
100
101    /// Merge segments into one, streaming postings/positions/store directly to files.
102    ///
103    /// If `trained` is provided, dense vectors use O(1) cluster merge when possible
104    /// (homogeneous IVF/ScaNN), otherwise rebuilds ANN from trained structures.
105    /// Without trained structures, only flat vectors are merged.
106    ///
107    /// Uses streaming writers so postings, positions, and store data flow directly
108    /// to files instead of buffering everything in memory. Only the term dictionary
109    /// (compact key+TermInfo entries) is buffered.
110    pub async fn merge<D: Directory + DirectoryWriter>(
111        &self,
112        dir: &D,
113        segments: &[SegmentReader],
114        new_segment_id: SegmentId,
115        trained: Option<&TrainedVectorStructures>,
116    ) -> Result<(SegmentMeta, MergeStats)> {
117        let mut stats = MergeStats::default();
118        let files = SegmentFiles::new(new_segment_id.0);
119
120        // === Two-stage merge to bound page cache pressure ===
121        //
122        // Stage 1: postings + store + fast_fields (concurrent)
123        //   Touches .term_dict, .postings, .positions, .store, .fast files.
124        //
125        // Stage 2: sparse + dense vectors (concurrent)
126        //   Touches .sparse, .vectors files.
127        //
128        // Running all phases concurrently caused OOM on large merges because
129        // mmap'd source files from all 16+ segments compete for page cache
130        // simultaneously (200+ GB of mmap'd data for BMP grids alone).
131        // Two stages halve the concurrent working set.
132        let merge_start = std::time::Instant::now();
133
134        // ── Stage 1: text + store + fast fields ─────────────────────────
135        let postings_fut = async {
136            let mut postings_writer =
137                OffsetWriter::new(dir.streaming_writer_cold(&files.postings).await?);
138            let mut positions_writer =
139                OffsetWriter::new(dir.streaming_writer_cold(&files.positions).await?);
140            let mut term_dict_writer =
141                OffsetWriter::new(dir.streaming_writer_cold(&files.term_dict).await?);
142
143            let terms_processed = self
144                .merge_postings(
145                    segments,
146                    &mut term_dict_writer,
147                    &mut postings_writer,
148                    &mut positions_writer,
149                )
150                .await?;
151
152            let postings_bytes = postings_writer.offset() as usize;
153            let term_dict_bytes = term_dict_writer.offset() as usize;
154            let positions_bytes = positions_writer.offset();
155
156            postings_writer.finish()?;
157            term_dict_writer.finish()?;
158            if positions_bytes > 0 {
159                positions_writer.finish()?;
160            } else {
161                drop(positions_writer);
162                let _ = dir.delete(&files.positions).await;
163            }
164            log::info!(
165                "[merge] postings done: {} terms, term_dict={}, postings={}, positions={}",
166                terms_processed,
167                format_bytes(term_dict_bytes),
168                format_bytes(postings_bytes),
169                format_bytes(positions_bytes as usize),
170            );
171            Ok::<(usize, usize, usize), crate::Error>((
172                terms_processed,
173                term_dict_bytes,
174                postings_bytes,
175            ))
176        };
177
178        let store_fut = async {
179            let mut store_writer =
180                OffsetWriter::new(dir.streaming_writer_cold(&files.store).await?);
181            let store_num_docs = self.merge_store(segments, &mut store_writer).await?;
182            let bytes = store_writer.offset() as usize;
183            store_writer.finish()?;
184            Ok::<(usize, u32), crate::Error>((bytes, store_num_docs))
185        };
186
187        let fast_fut = async { self.merge_fast_fields(dir, segments, &files).await };
188
189        let (postings_result, store_result, fast_bytes) =
190            tokio::try_join!(postings_fut, store_fut, fast_fut)?;
191
192        log::info!(
193            "[merge] stage 1 done in {:.1}s (postings + store + fast)",
194            merge_start.elapsed().as_secs_f64()
195        );
196
197        // ── Stage 2: sparse + dense vectors ─────────────────────────────
198        // Page cache from stage 1 files can now be evicted by the kernel
199        // as stage 2 accesses different mmap regions (.sparse, .vectors).
200        let sparse_fut = async { self.merge_sparse_vectors(dir, segments, &files).await };
201
202        let dense_fut = async {
203            self.merge_dense_vectors(dir, segments, &files, trained)
204                .await
205        };
206
207        let (sparse_bytes, vectors_bytes) = tokio::try_join!(sparse_fut, dense_fut)?;
208        let (store_bytes, store_num_docs) = store_result;
209        stats.terms_processed = postings_result.0;
210        stats.term_dict_bytes = postings_result.1;
211        stats.postings_bytes = postings_result.2;
212        stats.store_bytes = store_bytes;
213        stats.vectors_bytes = vectors_bytes;
214        stats.sparse_bytes = sparse_bytes;
215        stats.fast_bytes = fast_bytes;
216        log::info!(
217            "[merge] all phases done in {:.1}s: {}",
218            merge_start.elapsed().as_secs_f64(),
219            stats
220        );
221
222        // === Mandatory: merge field stats + write meta ===
223        let mut merged_field_stats: FxHashMap<u32, FieldStats> = FxHashMap::default();
224        for segment in segments {
225            for (&field_id, field_stats) in &segment.meta().field_stats {
226                let entry = merged_field_stats.entry(field_id).or_default();
227                entry.total_tokens += field_stats.total_tokens;
228                entry.doc_count += field_stats.doc_count;
229            }
230        }
231
232        let total_docs: u32 = segments
233            .iter()
234            .try_fold(0u32, |acc, s| acc.checked_add(s.num_docs()))
235            .ok_or_else(|| {
236                crate::Error::Internal(format!(
237                    "Total document count exceeds u32::MAX ({})",
238                    u32::MAX
239                ))
240            })?;
241
242        // Verify store doc count matches metadata — a mismatch here means
243        // some store blocks were lost (e.g., compression thread panic) or
244        // source segment metadata disagrees with its store.
245        if store_num_docs != total_docs {
246            log::error!(
247                "[merge] STORE/META MISMATCH: store has {} docs but metadata expects {}. \
248                 Per-segment: {:?}",
249                store_num_docs,
250                total_docs,
251                segments
252                    .iter()
253                    .map(|s| (
254                        format!("{:016x}", s.meta().id),
255                        s.num_docs(),
256                        s.store().num_docs()
257                    ))
258                    .collect::<Vec<_>>()
259            );
260            return Err(crate::Error::Io(std::io::Error::new(
261                std::io::ErrorKind::InvalidData,
262                format!(
263                    "Store/meta doc count mismatch: store={}, meta={}",
264                    store_num_docs, total_docs
265                ),
266            )));
267        }
268
269        let meta = SegmentMeta {
270            id: new_segment_id.0,
271            num_docs: total_docs,
272            field_stats: merged_field_stats,
273        };
274
275        dir.write(&files.meta, &meta.serialize()?).await?;
276
277        let label = if trained.is_some() {
278            "ANN merge"
279        } else {
280            "Merge"
281        };
282        log::info!("{} complete: {} docs, {}", label, total_docs, stats);
283
284        Ok((meta, stats))
285    }
286}
287
288/// Delete segment files from directory (all deletions run concurrently).
289pub async fn delete_segment<D: Directory + DirectoryWriter>(
290    dir: &D,
291    segment_id: SegmentId,
292) -> Result<()> {
293    let files = SegmentFiles::new(segment_id.0);
294    let _ = tokio::join!(
295        dir.delete(&files.term_dict),
296        dir.delete(&files.postings),
297        dir.delete(&files.store),
298        dir.delete(&files.meta),
299        dir.delete(&files.vectors),
300        dir.delete(&files.sparse),
301        dir.delete(&files.positions),
302        dir.delete(&files.fast),
303    );
304    Ok(())
305}