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::{FieldType, Schema};
21use crate::index::{ReorderConcurrencyGate, ReorderPriority};
22use crate::structures::SparseFormat;
23
24/// Compute per-segment doc ID offsets (each segment's docs start after the previous).
25///
26/// Returns an error if the total document count across segments exceeds `u32::MAX`.
27fn doc_offsets(segments: &[SegmentReader]) -> Result<Vec<u32>> {
28    let mut offsets = Vec::with_capacity(segments.len());
29    let mut acc = 0u32;
30    for seg in segments {
31        offsets.push(acc);
32        acc = acc.checked_add(seg.num_docs()).ok_or_else(|| {
33            crate::Error::Internal(format!(
34                "Total document count across segments exceeds u32::MAX ({})",
35                u32::MAX
36            ))
37        })?;
38    }
39    Ok(offsets)
40}
41
42/// Additive count stored in a `u32` field of the merged segment format.
43///
44/// Source segments are individually valid, so exceeding the limit is a
45/// property of this merge plan rather than source corruption.
46#[derive(Clone, Copy, Debug, Default)]
47struct MergeCapacity(u64);
48
49impl MergeCapacity {
50    #[inline]
51    fn add(&mut self, count: u64) -> Option<u64> {
52        self.0 = self.0.saturating_add(count);
53        (self.0 > u64::from(u32::MAX)).then_some(self.0)
54    }
55}
56
57fn field_capacity_error(
58    field_id: u32,
59    field_name: &str,
60    value_kind: &str,
61    count: u64,
62) -> crate::Error {
63    crate::Error::Schema(format!(
64        "merge would produce {count} {value_kind} for field {field_id} ('{field_name}'), \
65         exceeding the segment format limit {}; lower max_segment_docs for this \
66         multi-valued field",
67        u32::MAX,
68    ))
69}
70
71/// Statistics for merge operations
72#[derive(Debug, Clone, Default)]
73pub struct MergeStats {
74    /// Number of terms processed
75    pub terms_processed: usize,
76    /// Term dictionary output size
77    pub term_dict_bytes: usize,
78    /// Postings output size
79    pub postings_bytes: usize,
80    /// Store output size
81    pub store_bytes: usize,
82    /// Vector index output size
83    pub vectors_bytes: usize,
84    /// Sparse vector index output size
85    pub sparse_bytes: usize,
86    /// Whether merge-time BP reorder ran to full depth on every BMP field
87    /// (false = a pass hit its wall-clock budget; the segment is valid and
88    /// better-ordered, and the background optimizer deepens it later).
89    /// True when no BP ran (block-copy merges have nothing to deepen... they
90    /// are simply not reordered and tracked by the `reordered` flag instead).
91    pub bp_converged: bool,
92    /// Fast-field output size
93    pub fast_bytes: usize,
94}
95
96impl std::fmt::Display for MergeStats {
97    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
98        write!(
99            f,
100            "terms={}, term_dict={}, postings={}, store={}, dense_vectors={}, sparse_vectors={}, fast_fields={}",
101            self.terms_processed,
102            crate::format_bytes(self.term_dict_bytes as u64),
103            crate::format_bytes(self.postings_bytes as u64),
104            crate::format_bytes(self.store_bytes as u64),
105            crate::format_bytes(self.vectors_bytes as u64),
106            crate::format_bytes(self.sparse_bytes as u64),
107            crate::format_bytes(self.fast_bytes as u64),
108        )
109    }
110}
111
112// TrainedVectorStructures is defined in super::types (available on all platforms)
113pub use super::types::TrainedVectorStructures;
114
115/// Run a CPU/IO-heavy synchronous section, telling tokio to migrate this
116/// worker's task queue first (multi-thread runtimes only — `block_in_place`
117/// panics on current_thread, where we just run inline).
118pub(crate) fn block_in_place_if_multithread<R>(f: impl FnOnce() -> R) -> R {
119    if tokio::runtime::Handle::try_current()
120        .map(|h| h.runtime_flavor() == tokio::runtime::RuntimeFlavor::MultiThread)
121        .unwrap_or(false)
122    {
123        tokio::task::block_in_place(f)
124    } else {
125        f()
126    }
127}
128
129/// Append an exact-length temporary directory file to a segment output and
130/// remove it. Used by sparse skip tables so neither merge nor BP rewrite
131/// buffers a corpus-sized metadata section on heap.
132pub(crate) async fn append_and_delete_temp<D: DirectoryWriter>(
133    directory: &D,
134    path: &std::path::Path,
135    expected_bytes: u64,
136    writer: &mut OffsetWriter,
137) -> Result<()> {
138    use std::io::Write as _;
139
140    const COPY_CHUNK: u64 = 4 * 1024 * 1024;
141    let actual_bytes = directory.file_size(path).await?;
142    if actual_bytes != expected_bytes {
143        return Err(crate::Error::Corruption(format!(
144            "temporary sparse section {:?} has {} bytes, expected {}",
145            path, actual_bytes, expected_bytes,
146        )));
147    }
148    let mut offset = 0u64;
149    while offset < expected_bytes {
150        let end = (offset + COPY_CHUNK).min(expected_bytes);
151        let chunk = directory.read_range(path, offset..end).await?;
152        writer
153            .write_all(chunk.as_slice())
154            .map_err(crate::Error::Io)?;
155        offset = end;
156    }
157    if let Err(error) = directory.delete(path).await {
158        // The section is already complete in the output. This output-scoped
159        // scratch file is safe for the startup orphan sweep and must not
160        // invalidate an otherwise successful multi-hour merge.
161        log::warn!(
162            "[merge] failed to remove temporary sparse section {:?}: {}",
163            path,
164            error,
165        );
166    }
167    Ok(())
168}
169
170/// Segment merger - merges multiple segments into one
171pub struct SegmentMerger {
172    schema: Arc<Schema>,
173    /// Run BP reordering on BMP sparse fields while writing the merged blob
174    /// (instead of byte-level block stacking). The output segment is then
175    /// already ordered, so the standalone reorder pass is unnecessary.
176    reorder_bmp: bool,
177    /// Bounded rayon pool for merge-time BP. `None` = global pool (tests);
178    /// the SegmentManager always passes its background pool so BP cannot
179    /// starve query scoring.
180    background_pool: Option<Arc<rayon::ThreadPool>>,
181    /// Granularity for merge-time BP. `Auto` by default; the SegmentManager
182    /// forces `Records` when any merge source is an unconverged partial
183    /// reorder.
184    granularity: crate::segment::reorder::BpGranularity,
185    /// Budget for merge-time BP. Default unbudgeted; the SegmentManager
186    /// passes the index's `merge_bp_time_budget` so huge merges stop holding
187    /// a merge slot for the full BP depth — a truncated pass is marked
188    /// `bp_converged = false` and the background optimizer deepens it.
189    bp_budget: crate::segment::BpBudget,
190    /// Memory budget for the BP forward index during merge-time reorder.
191    bp_memory_budget: usize,
192    /// Shared whole-pass concurrency limit. Tests and low-level callers may
193    /// omit it; SegmentManager always supplies the application-wide gate.
194    reorder_permits: Option<Arc<ReorderConcurrencyGate>>,
195    /// Automatic merges are background work. An explicit force merge holds a
196    /// foreground guard and bypasses the background pause for its BP fields.
197    reorder_priority: ReorderPriority,
198}
199
200impl SegmentMerger {
201    pub fn new(schema: Arc<Schema>) -> Self {
202        Self {
203            schema,
204            reorder_bmp: false,
205            background_pool: None,
206            granularity: crate::segment::reorder::BpGranularity::Auto,
207            bp_budget: crate::segment::BpBudget::full(),
208            bp_memory_budget: crate::segment::reorder::DEFAULT_MEMORY_BUDGET,
209            reorder_permits: None,
210            reorder_priority: ReorderPriority::Background,
211        }
212    }
213
214    /// Enable BP reordering of BMP fields during the merge (see `reorder_bmp`).
215    pub fn with_bmp_reorder(mut self, reorder: bool) -> Self {
216        self.reorder_bmp = reorder;
217        self
218    }
219
220    /// Run merge-time BP on this bounded pool instead of the global one.
221    pub fn with_background_pool(mut self, pool: Option<Arc<rayon::ThreadPool>>) -> Self {
222        self.background_pool = pool;
223        self
224    }
225
226    /// Set merge-time BP granularity (see `granularity`).
227    pub fn with_granularity(mut self, granularity: crate::segment::reorder::BpGranularity) -> Self {
228        self.granularity = granularity;
229        self
230    }
231
232    /// Bound merge-time BP wall clock (see `bp_budget`).
233    pub fn with_bp_budget(mut self, budget: crate::segment::BpBudget) -> Self {
234        self.bp_budget = budget;
235        self
236    }
237
238    /// Memory budget for the BP forward index (see `bp_memory_budget`).
239    pub fn with_bp_memory_budget(mut self, bytes: usize) -> Self {
240        self.bp_memory_budget = bytes;
241        self
242    }
243
244    /// Share the application-wide whole-segment reorder gate.
245    pub fn with_reorder_permits(mut self, permits: Arc<ReorderConcurrencyGate>) -> Self {
246        self.reorder_permits = Some(permits);
247        self
248    }
249
250    pub(crate) fn with_reorder_priority(mut self, priority: ReorderPriority) -> Self {
251        self.reorder_priority = priority;
252        self
253    }
254
255    /// Reject additive per-field counts that the on-disk formats cannot
256    /// represent. All inputs are already-open metadata views; no vector,
257    /// posting, or document payload is read here.
258    fn validate_merge_capacities(&self, segments: &[SegmentReader]) -> Result<()> {
259        // MaxScore skip entries share one u32-addressed section across fields.
260        let mut maxscore_skip_entries = MergeCapacity::default();
261
262        for (field, entry) in self.schema.fields() {
263            match entry.field_type {
264                FieldType::DenseVector | FieldType::BinaryDenseVector => {
265                    let mut vectors = MergeCapacity::default();
266                    for segment in segments {
267                        let Some(flat) = segment.flat_vectors().get(&field.0) else {
268                            continue;
269                        };
270                        if let Some(total) = vectors.add(flat.num_vectors as u64) {
271                            let value_kind = if entry.field_type == FieldType::BinaryDenseVector {
272                                "binary vectors"
273                            } else {
274                                "dense vectors"
275                            };
276                            return Err(field_capacity_error(
277                                field.0,
278                                &entry.name,
279                                value_kind,
280                                total,
281                            ));
282                        }
283                    }
284                }
285                FieldType::SparseVector => {
286                    let format = entry
287                        .sparse_vector_config
288                        .as_ref()
289                        .map(|config| config.format)
290                        .unwrap_or_default();
291                    match format {
292                        SparseFormat::Bmp => {
293                            let mut vectors = MergeCapacity::default();
294                            let mut blocks = MergeCapacity::default();
295                            let mut real_slots = MergeCapacity::default();
296                            let mut virtual_slots = MergeCapacity::default();
297
298                            for segment in segments {
299                                let Some(index) = segment.bmp_indexes().get(&field.0) else {
300                                    continue;
301                                };
302                                for (capacity, count, value_kind) in [
303                                    (&mut vectors, u64::from(index.total_vectors), "BMP vectors"),
304                                    (&mut blocks, u64::from(index.num_blocks), "BMP blocks"),
305                                    (
306                                        &mut real_slots,
307                                        u64::from(index.num_real_docs()),
308                                        "BMP real vector slots",
309                                    ),
310                                    (
311                                        &mut virtual_slots,
312                                        u64::from(index.num_virtual_docs),
313                                        "BMP padded virtual slots",
314                                    ),
315                                ] {
316                                    if let Some(total) = capacity.add(count) {
317                                        return Err(field_capacity_error(
318                                            field.0,
319                                            &entry.name,
320                                            value_kind,
321                                            total,
322                                        ));
323                                    }
324                                }
325                            }
326                        }
327                        SparseFormat::MaxScore => {
328                            let mut vectors = MergeCapacity::default();
329                            let mut dimensions: FxHashMap<u32, (MergeCapacity, MergeCapacity)> =
330                                FxHashMap::default();
331
332                            for segment in segments {
333                                let Some(index) = segment.sparse_indexes().get(&field.0) else {
334                                    continue;
335                                };
336                                if let Some(total) = vectors.add(u64::from(index.total_vectors)) {
337                                    return Err(field_capacity_error(
338                                        field.0,
339                                        &entry.name,
340                                        "MaxScore vectors",
341                                        total,
342                                    ));
343                                }
344
345                                for (dimension, doc_count, block_count) in index.dimension_counts()
346                                {
347                                    let (docs, blocks) = dimensions.entry(dimension).or_default();
348                                    if let Some(total) = docs.add(u64::from(doc_count)) {
349                                        return Err(field_capacity_error(
350                                            field.0,
351                                            &entry.name,
352                                            &format!("MaxScore postings for dimension {dimension}"),
353                                            total,
354                                        ));
355                                    }
356                                    if let Some(total) = blocks.add(u64::from(block_count)) {
357                                        return Err(field_capacity_error(
358                                            field.0,
359                                            &entry.name,
360                                            &format!("MaxScore blocks for dimension {dimension}"),
361                                            total,
362                                        ));
363                                    }
364                                    if let Some(total) =
365                                        maxscore_skip_entries.add(u64::from(block_count))
366                                    {
367                                        return Err(crate::Error::Schema(format!(
368                                            "merge would produce {total} MaxScore skip entries \
369                                             across sparse fields, exceeding the segment format \
370                                             limit {}; lower max_segment_docs for multi-valued \
371                                             sparse fields",
372                                            u32::MAX,
373                                        )));
374                                    }
375                                }
376                            }
377                        }
378                    }
379                }
380                _ => {}
381            }
382        }
383        Ok(())
384    }
385
386    /// Merge segments into one, streaming postings/positions/store directly to files.
387    ///
388    /// If `trained` is provided, dense vectors use O(1) cluster merge when possible
389    /// (compatible IVF-PQ), otherwise rebuilds ANN from global artifacts.
390    /// Without trained structures, only flat vectors are merged.
391    ///
392    /// Uses streaming writers so postings, positions, and store data flow directly
393    /// to files instead of buffering everything in memory. Only the term dictionary
394    /// (compact key+TermInfo entries) is buffered.
395    pub async fn merge<D: Directory + DirectoryWriter>(
396        &self,
397        dir: &D,
398        segments: &[SegmentReader],
399        new_segment_id: SegmentId,
400        trained: Option<&TrainedVectorStructures>,
401    ) -> Result<(SegmentMeta, MergeStats)> {
402        // Reject an unrepresentable merge before creating any output files.
403        // The previous late check left a complete orphan output behind after
404        // doing all expensive phases.
405        let total_docs: u32 = segments
406            .iter()
407            .try_fold(0u32, |acc, segment| acc.checked_add(segment.num_docs()))
408            .ok_or_else(|| {
409                crate::Error::Internal(format!(
410                    "Total document count exceeds u32::MAX ({})",
411                    u32::MAX
412                ))
413            })?;
414
415        self.validate_merge_capacities(segments)?;
416
417        let mut stats = MergeStats::default();
418        let files = SegmentFiles::new(new_segment_id.0);
419
420        // === Two-stage merge to bound page cache pressure ===
421        //
422        // Stage 1: postings + store + fast_fields (concurrent)
423        //   Touches .term_dict, .postings, .positions, .store, .fast files.
424        //
425        // Stage 2: sparse + dense vectors. Block-copy sparse work runs with
426        // dense vectors; BP sparse work runs first to bound peak memory.
427        //   Touches .sparse, .vectors files.
428        //
429        // Running all phases concurrently caused OOM on large merges because
430        // mmap'd source files from all 16+ segments compete for page cache
431        // simultaneously (200+ GB of mmap'd data for BMP grids alone).
432        // Two stages halve the concurrent working set.
433        let merge_start = std::time::Instant::now();
434
435        // ── Stage 1: text + store + fast fields ─────────────────────────
436        let postings_fut = async {
437            let mut postings_writer =
438                OffsetWriter::new(dir.streaming_writer_cold(&files.postings).await?);
439            let mut positions_writer =
440                OffsetWriter::new(dir.streaming_writer_cold(&files.positions).await?);
441            let mut term_dict_writer =
442                OffsetWriter::new(dir.streaming_writer_cold(&files.term_dict).await?);
443
444            let terms_processed = self
445                .merge_postings(
446                    segments,
447                    &mut term_dict_writer,
448                    &mut postings_writer,
449                    &mut positions_writer,
450                )
451                .await?;
452
453            let postings_bytes = postings_writer.offset() as usize;
454            let term_dict_bytes = term_dict_writer.offset() as usize;
455            let positions_bytes = positions_writer.offset();
456
457            postings_writer.finish()?;
458            term_dict_writer.finish()?;
459            if positions_bytes > 0 {
460                positions_writer.finish()?;
461            } else {
462                drop(positions_writer);
463                let _ = dir.delete(&files.positions).await;
464            }
465            log::info!(
466                "[merge] postings done: {} terms, term_dict={}, postings={}, positions={}",
467                terms_processed,
468                crate::format_bytes(term_dict_bytes as u64),
469                crate::format_bytes(postings_bytes as u64),
470                crate::format_bytes(positions_bytes),
471            );
472            Ok::<(usize, usize, usize), crate::Error>((
473                terms_processed,
474                term_dict_bytes,
475                postings_bytes,
476            ))
477        };
478
479        let store_fut = async {
480            let mut store_writer =
481                OffsetWriter::new(dir.streaming_writer_cold(&files.store).await?);
482            let store_num_docs = self.merge_store(segments, &mut store_writer).await?;
483            let bytes = store_writer.offset() as usize;
484            store_writer.finish()?;
485            Ok::<(usize, u32), crate::Error>((bytes, store_num_docs))
486        };
487
488        let fast_fut = async { self.merge_fast_fields(dir, segments, &files).await };
489
490        let (postings_result, store_result, fast_bytes) =
491            tokio::try_join!(postings_fut, store_fut, fast_fut)?;
492
493        log::info!(
494            "[merge] stage 1 done in {:.1}s (postings + store + fast)",
495            merge_start.elapsed().as_secs_f64()
496        );
497
498        // ── Stage 2: sparse + dense vectors ─────────────────────────────
499        // Page cache from stage 1 files can now be evicted by the kernel
500        // as stage 2 accesses different mmap regions (.sparse, .vectors).
501        let sparse_fut = async { self.merge_sparse_vectors(dir, segments, &files).await };
502
503        let dense_fut = async {
504            self.merge_dense_vectors(dir, segments, &files, trained, AnnWriteMode::Copy)
505                .await
506        };
507
508        // Merge-time BP constructs a potentially budget-sized forward index.
509        // Do not overlap that allocation and its heavy source-file scan with
510        // an ANN rebuild. Block-copy sparse merges remain concurrent with ANN.
511        let ((sparse_bytes, bp_converged), vectors_bytes) = if self.reorder_bmp {
512            let sparse = sparse_fut.await?;
513            let dense = dense_fut.await?;
514            (sparse, dense)
515        } else {
516            tokio::try_join!(sparse_fut, dense_fut)?
517        };
518        let (store_bytes, store_num_docs) = store_result;
519        stats.terms_processed = postings_result.0;
520        stats.term_dict_bytes = postings_result.1;
521        stats.postings_bytes = postings_result.2;
522        stats.store_bytes = store_bytes;
523        stats.vectors_bytes = vectors_bytes;
524        stats.sparse_bytes = sparse_bytes;
525        stats.bp_converged = bp_converged;
526        stats.fast_bytes = fast_bytes;
527        log::info!(
528            "[merge] all phases done in {:.1}s: {}",
529            merge_start.elapsed().as_secs_f64(),
530            stats
531        );
532
533        // === Mandatory: merge field stats + write meta ===
534        let mut merged_field_stats: FxHashMap<u32, FieldStats> = FxHashMap::default();
535        for segment in segments {
536            for (&field_id, field_stats) in &segment.meta().field_stats {
537                let entry = merged_field_stats.entry(field_id).or_default();
538                entry.total_tokens = entry
539                    .total_tokens
540                    .checked_add(field_stats.total_tokens)
541                    .ok_or_else(|| {
542                        crate::Error::Corruption(format!(
543                            "field {} total-token count overflow while merging",
544                            field_id
545                        ))
546                    })?;
547                entry.doc_count = entry
548                    .doc_count
549                    .checked_add(field_stats.doc_count)
550                    .ok_or_else(|| {
551                        crate::Error::Corruption(format!(
552                            "field {} document count overflow while merging",
553                            field_id
554                        ))
555                    })?;
556            }
557        }
558
559        // Verify store doc count matches metadata — a mismatch here means
560        // some store blocks were lost (e.g., compression thread panic) or
561        // source segment metadata disagrees with its store.
562        if store_num_docs != total_docs {
563            log::error!(
564                "[merge] STORE/META MISMATCH: store has {} docs but metadata expects {}. \
565                 Per-segment: {:?}",
566                store_num_docs,
567                total_docs,
568                segments
569                    .iter()
570                    .map(|s| (
571                        format!("{:016x}", s.meta().id),
572                        s.num_docs(),
573                        s.store().num_docs()
574                    ))
575                    .collect::<Vec<_>>()
576            );
577            return Err(crate::Error::Io(std::io::Error::new(
578                std::io::ErrorKind::InvalidData,
579                format!(
580                    "Store/meta doc count mismatch: store={}, meta={}",
581                    store_num_docs, total_docs
582                ),
583            )));
584        }
585
586        let meta = SegmentMeta {
587            id: new_segment_id.0,
588            num_docs: total_docs,
589            field_stats: merged_field_stats,
590        };
591
592        // Durable: replace_segments deletes the fsynced source segments right
593        // after publishing this output, so a non-durable .meta could be the
594        // only copy of the merged documents across a power failure.
595        dir.write_durable(&files.meta, &meta.serialize()?).await?;
596
597        let label = if trained.is_some() {
598            "ANN merge"
599        } else {
600            "Merge"
601        };
602        log::info!("{} complete: {} docs, {}", label, total_docs, stats);
603
604        Ok((meta, stats))
605    }
606}
607
608/// Delete segment files from directory (all deletions run concurrently).
609pub async fn delete_segment<D: Directory + DirectoryWriter>(
610    dir: &D,
611    segment_id: SegmentId,
612) -> Result<()> {
613    let files = SegmentFiles::new(segment_id.0);
614    let paths = files.lifecycle_paths();
615    let results = futures::future::join_all(paths.iter().map(|path| dir.delete(path))).await;
616
617    // Missing files are expected for optional components and idempotent
618    // retries. Any other failure must be surfaced so cleanup is not falsely
619    // reported as successful; a later orphan sweep can retry remaining files.
620    for result in results {
621        if let Err(error) = result
622            && error.kind() != std::io::ErrorKind::NotFound
623        {
624            return Err(crate::Error::Io(error));
625        }
626    }
627    Ok(())
628}
629
630#[cfg(test)]
631mod capacity_tests {
632    use super::{MergeCapacity, field_capacity_error};
633
634    #[test]
635    fn merge_capacity_accepts_the_exact_u32_boundary() {
636        let mut capacity = MergeCapacity::default();
637        assert_eq!(capacity.add(u64::from(u32::MAX) - 7), None);
638        assert_eq!(capacity.add(7), None);
639    }
640
641    #[test]
642    fn merge_capacity_rejects_the_first_value_beyond_u32() {
643        let mut capacity = MergeCapacity::default();
644        assert_eq!(capacity.add(u64::from(u32::MAX)), None);
645        assert_eq!(capacity.add(1), Some(u64::from(u32::MAX) + 1));
646    }
647
648    #[test]
649    fn merge_capacity_failure_is_not_source_corruption() {
650        let error = field_capacity_error(
651            7,
652            "body_embedding",
653            "dense vectors",
654            u64::from(u32::MAX) + 1,
655        );
656        assert!(matches!(error, crate::Error::Schema(_)));
657        assert!(error.to_string().contains("lower max_segment_docs"));
658    }
659}