1mod dense;
4mod fast_fields;
5mod postings;
6mod sparse;
7mod store;
8
9pub(crate) use dense::AnnWriteMode;
10
11use std::sync::Arc;
12use std::sync::atomic::{AtomicBool, Ordering};
13
14use rustc_hash::FxHashMap;
15
16use super::OffsetWriter;
17use super::reader::SegmentReader;
18use super::types::{FieldStats, SegmentFiles, SegmentId, SegmentMeta};
19use crate::Result;
20use crate::directories::{Directory, DirectoryWriter};
21use crate::dsl::{FieldType, Schema};
22use crate::index::{ReorderConcurrencyGate, ReorderPriority};
23use crate::structures::SparseFormat;
24
25fn doc_offsets(segments: &[SegmentReader]) -> Result<Vec<u32>> {
29 let mut offsets = Vec::with_capacity(segments.len());
30 let mut acc = 0u32;
31 for seg in segments {
32 offsets.push(acc);
33 acc = acc.checked_add(seg.num_docs()).ok_or_else(|| {
34 crate::Error::Internal(format!(
35 "Total document count across segments exceeds u32::MAX ({})",
36 u32::MAX
37 ))
38 })?;
39 }
40 Ok(offsets)
41}
42
43#[derive(Clone, Copy, Debug, Default)]
48struct MergeCapacity(u64);
49
50impl MergeCapacity {
51 #[inline]
52 fn add(&mut self, count: u64) -> Option<u64> {
53 self.0 = self.0.saturating_add(count);
54 (self.0 > u64::from(u32::MAX)).then_some(self.0)
55 }
56}
57
58fn field_capacity_error(
59 field_id: u32,
60 field_name: &str,
61 value_kind: &str,
62 count: u64,
63) -> crate::Error {
64 crate::Error::Schema(format!(
65 "merge would produce {count} {value_kind} for field {field_id} ('{field_name}'), \
66 exceeding the segment format limit {}; lower max_segment_docs for this \
67 multi-valued field",
68 u32::MAX,
69 ))
70}
71
72#[derive(Debug, Clone, Default)]
74pub struct MergeStats {
75 pub terms_processed: usize,
77 pub term_dict_bytes: usize,
79 pub postings_bytes: usize,
81 pub store_bytes: usize,
83 pub vectors_bytes: usize,
85 pub sparse_bytes: usize,
87 pub bp_converged: bool,
93 pub fast_bytes: usize,
95}
96
97impl std::fmt::Display for MergeStats {
98 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
99 write!(
100 f,
101 "terms={}, term_dict={}, postings={}, store={}, dense_vectors={}, sparse_vectors={}, fast_fields={}",
102 self.terms_processed,
103 crate::format_bytes(self.term_dict_bytes as u64),
104 crate::format_bytes(self.postings_bytes as u64),
105 crate::format_bytes(self.store_bytes as u64),
106 crate::format_bytes(self.vectors_bytes as u64),
107 crate::format_bytes(self.sparse_bytes as u64),
108 crate::format_bytes(self.fast_bytes as u64),
109 )
110 }
111}
112
113pub use super::types::TrainedVectorStructures;
115
116pub(crate) fn block_in_place_if_multithread<R>(f: impl FnOnce() -> R) -> R {
120 if tokio::runtime::Handle::try_current()
121 .map(|h| h.runtime_flavor() == tokio::runtime::RuntimeFlavor::MultiThread)
122 .unwrap_or(false)
123 {
124 tokio::task::block_in_place(f)
125 } else {
126 f()
127 }
128}
129
130pub(crate) async fn append_and_delete_temp<D: DirectoryWriter>(
134 directory: &D,
135 path: &std::path::Path,
136 expected_bytes: u64,
137 writer: &mut OffsetWriter,
138) -> Result<()> {
139 use std::io::Write as _;
140
141 const COPY_CHUNK: u64 = 4 * 1024 * 1024;
142 let actual_bytes = directory.file_size(path).await?;
143 if actual_bytes != expected_bytes {
144 return Err(crate::Error::Corruption(format!(
145 "temporary sparse section {:?} has {} bytes, expected {}",
146 path, actual_bytes, expected_bytes,
147 )));
148 }
149 let mut offset = 0u64;
150 while offset < expected_bytes {
151 let end = (offset + COPY_CHUNK).min(expected_bytes);
152 let chunk = directory.read_range(path, offset..end).await?;
153 writer
154 .write_all(chunk.as_slice())
155 .map_err(crate::Error::Io)?;
156 offset = end;
157 }
158 if let Err(error) = directory.delete(path).await {
159 log::warn!(
163 "[merge] failed to remove temporary sparse section {:?}: {}",
164 path,
165 error,
166 );
167 }
168 Ok(())
169}
170
171pub struct SegmentMerger {
173 schema: Arc<Schema>,
174 reorder_bmp: bool,
178 background_pool: Option<Arc<rayon::ThreadPool>>,
182 granularity: crate::segment::reorder::BpGranularity,
186 bp_budget: crate::segment::BpBudget,
191 cancellation: Option<Arc<AtomicBool>>,
194 bp_memory_budget: usize,
196 reorder_permits: Option<Arc<ReorderConcurrencyGate>>,
199 reorder_priority: ReorderPriority,
202}
203
204impl SegmentMerger {
205 pub fn new(schema: Arc<Schema>) -> Self {
206 Self {
207 schema,
208 reorder_bmp: false,
209 background_pool: None,
210 granularity: crate::segment::reorder::BpGranularity::Auto,
211 bp_budget: crate::segment::BpBudget::full(),
212 cancellation: None,
213 bp_memory_budget: crate::segment::reorder::DEFAULT_MEMORY_BUDGET,
214 reorder_permits: None,
215 reorder_priority: ReorderPriority::AutomaticMerge,
216 }
217 }
218
219 pub fn with_bmp_reorder(mut self, reorder: bool) -> Self {
221 self.reorder_bmp = reorder;
222 self
223 }
224
225 pub fn with_background_pool(mut self, pool: Option<Arc<rayon::ThreadPool>>) -> Self {
227 self.background_pool = pool;
228 self
229 }
230
231 pub fn with_granularity(mut self, granularity: crate::segment::reorder::BpGranularity) -> Self {
233 self.granularity = granularity;
234 self
235 }
236
237 pub fn with_bp_budget(mut self, budget: crate::segment::BpBudget) -> Self {
239 self.bp_budget = budget;
240 self
241 }
242
243 pub(crate) fn with_cancellation(mut self, cancellation: Arc<AtomicBool>) -> Self {
244 self.cancellation = Some(cancellation);
245 self
246 }
247
248 pub fn with_bp_memory_budget(mut self, bytes: usize) -> Self {
250 self.bp_memory_budget = bytes;
251 self
252 }
253
254 pub fn with_reorder_permits(mut self, permits: Arc<ReorderConcurrencyGate>) -> Self {
256 self.reorder_permits = Some(permits);
257 self
258 }
259
260 pub(crate) fn with_reorder_priority(mut self, priority: ReorderPriority) -> Self {
261 self.reorder_priority = priority;
262 self
263 }
264
265 pub(super) fn ensure_not_cancelled(&self) -> Result<()> {
266 if self
267 .cancellation
268 .as_ref()
269 .is_some_and(|cancelled| cancelled.load(Ordering::Acquire))
270 {
271 Err(crate::Error::IndexClosed)
272 } else {
273 Ok(())
274 }
275 }
276
277 fn validate_merge_capacities(&self, segments: &[SegmentReader]) -> Result<()> {
281 let mut maxscore_skip_entries = MergeCapacity::default();
283
284 for (field, entry) in self.schema.fields() {
285 match entry.field_type {
286 FieldType::DenseVector | FieldType::BinaryDenseVector => {
287 let mut vectors = MergeCapacity::default();
288 for segment in segments {
289 let Some(flat) = segment.flat_vectors().get(&field.0) else {
290 continue;
291 };
292 if let Some(total) = vectors.add(flat.num_vectors as u64) {
293 let value_kind = if entry.field_type == FieldType::BinaryDenseVector {
294 "binary vectors"
295 } else {
296 "dense vectors"
297 };
298 return Err(field_capacity_error(
299 field.0,
300 &entry.name,
301 value_kind,
302 total,
303 ));
304 }
305 }
306 }
307 FieldType::SparseVector => {
308 let format = entry
309 .sparse_vector_config
310 .as_ref()
311 .map(|config| config.format)
312 .unwrap_or_default();
313 match format {
314 SparseFormat::Bmp => {
315 let mut vectors = MergeCapacity::default();
316 let mut blocks = MergeCapacity::default();
317 let mut real_slots = MergeCapacity::default();
318 let mut virtual_slots = MergeCapacity::default();
319
320 for segment in segments {
321 let Some(index) = segment.bmp_indexes().get(&field.0) else {
322 continue;
323 };
324 for (capacity, count, value_kind) in [
325 (&mut vectors, u64::from(index.total_vectors), "BMP vectors"),
326 (&mut blocks, u64::from(index.num_blocks), "BMP blocks"),
327 (
328 &mut real_slots,
329 u64::from(index.num_real_docs()),
330 "BMP real vector slots",
331 ),
332 (
333 &mut virtual_slots,
334 u64::from(index.num_virtual_docs),
335 "BMP padded virtual slots",
336 ),
337 ] {
338 if let Some(total) = capacity.add(count) {
339 return Err(field_capacity_error(
340 field.0,
341 &entry.name,
342 value_kind,
343 total,
344 ));
345 }
346 }
347 }
348 }
349 SparseFormat::MaxScore => {
350 let mut vectors = MergeCapacity::default();
351 let mut dimensions: FxHashMap<u32, (MergeCapacity, MergeCapacity)> =
352 FxHashMap::default();
353
354 for segment in segments {
355 let Some(index) = segment.sparse_indexes().get(&field.0) else {
356 continue;
357 };
358 if let Some(total) = vectors.add(u64::from(index.total_vectors)) {
359 return Err(field_capacity_error(
360 field.0,
361 &entry.name,
362 "MaxScore vectors",
363 total,
364 ));
365 }
366
367 for (dimension, doc_count, block_count) in index.dimension_counts()
368 {
369 let (docs, blocks) = dimensions.entry(dimension).or_default();
370 if let Some(total) = docs.add(u64::from(doc_count)) {
371 return Err(field_capacity_error(
372 field.0,
373 &entry.name,
374 &format!("MaxScore postings for dimension {dimension}"),
375 total,
376 ));
377 }
378 if let Some(total) = blocks.add(u64::from(block_count)) {
379 return Err(field_capacity_error(
380 field.0,
381 &entry.name,
382 &format!("MaxScore blocks for dimension {dimension}"),
383 total,
384 ));
385 }
386 if let Some(total) =
387 maxscore_skip_entries.add(u64::from(block_count))
388 {
389 return Err(crate::Error::Schema(format!(
390 "merge would produce {total} MaxScore skip entries \
391 across sparse fields, exceeding the segment format \
392 limit {}; lower max_segment_docs for multi-valued \
393 sparse fields",
394 u32::MAX,
395 )));
396 }
397 }
398 }
399 }
400 }
401 }
402 _ => {}
403 }
404 }
405 Ok(())
406 }
407
408 pub async fn merge<D: Directory + DirectoryWriter>(
418 &self,
419 dir: &D,
420 segments: &[SegmentReader],
421 new_segment_id: SegmentId,
422 trained: Option<&TrainedVectorStructures>,
423 ) -> Result<(SegmentMeta, MergeStats)> {
424 self.ensure_not_cancelled()?;
425 let total_docs: u32 = segments
429 .iter()
430 .try_fold(0u32, |acc, segment| acc.checked_add(segment.num_docs()))
431 .ok_or_else(|| {
432 crate::Error::Internal(format!(
433 "Total document count exceeds u32::MAX ({})",
434 u32::MAX
435 ))
436 })?;
437
438 self.validate_merge_capacities(segments)?;
439
440 let mut stats = MergeStats::default();
441 let files = SegmentFiles::new(new_segment_id.0);
442
443 let merge_start = std::time::Instant::now();
457
458 let postings_fut = async {
460 let mut postings_writer =
461 OffsetWriter::new(dir.streaming_writer_cold(&files.postings).await?);
462 let mut positions_writer =
463 OffsetWriter::new(dir.streaming_writer_cold(&files.positions).await?);
464 let mut term_dict_writer =
465 OffsetWriter::new(dir.streaming_writer_cold(&files.term_dict).await?);
466
467 let terms_processed = self
468 .merge_postings(
469 segments,
470 &mut term_dict_writer,
471 &mut postings_writer,
472 &mut positions_writer,
473 )
474 .await?;
475
476 let postings_bytes = postings_writer.offset() as usize;
477 let term_dict_bytes = term_dict_writer.offset() as usize;
478 let positions_bytes = positions_writer.offset();
479
480 postings_writer.finish()?;
481 term_dict_writer.finish()?;
482 if positions_bytes > 0 {
483 positions_writer.finish()?;
484 } else {
485 drop(positions_writer);
486 let _ = dir.delete(&files.positions).await;
487 }
488 log::info!(
489 "[merge] postings done: {} terms, term_dict={}, postings={}, positions={}",
490 terms_processed,
491 crate::format_bytes(term_dict_bytes as u64),
492 crate::format_bytes(postings_bytes as u64),
493 crate::format_bytes(positions_bytes),
494 );
495 Ok::<(usize, usize, usize), crate::Error>((
496 terms_processed,
497 term_dict_bytes,
498 postings_bytes,
499 ))
500 };
501
502 let store_fut = async {
503 let mut store_writer =
504 OffsetWriter::new(dir.streaming_writer_cold(&files.store).await?);
505 let store_num_docs = self.merge_store(segments, &mut store_writer).await?;
506 let bytes = store_writer.offset() as usize;
507 store_writer.finish()?;
508 Ok::<(usize, u32), crate::Error>((bytes, store_num_docs))
509 };
510
511 let fast_fut = async { self.merge_fast_fields(dir, segments, &files).await };
512
513 let (postings_result, store_result, fast_bytes) =
514 tokio::try_join!(postings_fut, store_fut, fast_fut)?;
515 self.ensure_not_cancelled()?;
516
517 log::info!(
518 "[merge] stage 1 done in {:.1}s (postings + store + fast)",
519 merge_start.elapsed().as_secs_f64()
520 );
521
522 let sparse_fut = async { self.merge_sparse_vectors(dir, segments, &files).await };
526
527 let dense_fut = async {
528 self.merge_dense_vectors(dir, segments, &files, trained, AnnWriteMode::Copy)
529 .await
530 };
531
532 let ((sparse_bytes, bp_converged), vectors_bytes) = if self.reorder_bmp {
536 let sparse = sparse_fut.await?;
537 let dense = dense_fut.await?;
538 (sparse, dense)
539 } else {
540 tokio::try_join!(sparse_fut, dense_fut)?
541 };
542 self.ensure_not_cancelled()?;
543 let (store_bytes, store_num_docs) = store_result;
544 stats.terms_processed = postings_result.0;
545 stats.term_dict_bytes = postings_result.1;
546 stats.postings_bytes = postings_result.2;
547 stats.store_bytes = store_bytes;
548 stats.vectors_bytes = vectors_bytes;
549 stats.sparse_bytes = sparse_bytes;
550 stats.bp_converged = bp_converged;
551 stats.fast_bytes = fast_bytes;
552 log::info!(
553 "[merge] all phases done in {:.1}s: {}",
554 merge_start.elapsed().as_secs_f64(),
555 stats
556 );
557
558 self.ensure_not_cancelled()?;
560 let mut merged_field_stats: FxHashMap<u32, FieldStats> = FxHashMap::default();
561 for segment in segments {
562 for (&field_id, field_stats) in &segment.meta().field_stats {
563 let entry = merged_field_stats.entry(field_id).or_default();
564 entry.total_tokens = entry
565 .total_tokens
566 .checked_add(field_stats.total_tokens)
567 .ok_or_else(|| {
568 crate::Error::Corruption(format!(
569 "field {} total-token count overflow while merging",
570 field_id
571 ))
572 })?;
573 entry.doc_count = entry
574 .doc_count
575 .checked_add(field_stats.doc_count)
576 .ok_or_else(|| {
577 crate::Error::Corruption(format!(
578 "field {} document count overflow while merging",
579 field_id
580 ))
581 })?;
582 }
583 }
584
585 if store_num_docs != total_docs {
589 log::error!(
590 "[merge] STORE/META MISMATCH: store has {} docs but metadata expects {}. \
591 Per-segment: {:?}",
592 store_num_docs,
593 total_docs,
594 segments
595 .iter()
596 .map(|s| (
597 format!("{:016x}", s.meta().id),
598 s.num_docs(),
599 s.store().num_docs()
600 ))
601 .collect::<Vec<_>>()
602 );
603 return Err(crate::Error::Io(std::io::Error::new(
604 std::io::ErrorKind::InvalidData,
605 format!(
606 "Store/meta doc count mismatch: store={}, meta={}",
607 store_num_docs, total_docs
608 ),
609 )));
610 }
611
612 let meta = SegmentMeta {
613 id: new_segment_id.0,
614 num_docs: total_docs,
615 field_stats: merged_field_stats,
616 };
617
618 dir.write_durable(&files.meta, &meta.serialize()?).await?;
622
623 let label = if trained.is_some() {
624 "ANN merge"
625 } else {
626 "Merge"
627 };
628 log::info!("{} complete: {} docs, {}", label, total_docs, stats);
629
630 Ok((meta, stats))
631 }
632}
633
634pub async fn delete_segment<D: Directory + DirectoryWriter>(
636 dir: &D,
637 segment_id: SegmentId,
638) -> Result<()> {
639 let files = SegmentFiles::new(segment_id.0);
640 let paths = files.lifecycle_paths();
641 let results = futures::future::join_all(paths.iter().map(|path| dir.delete(path))).await;
642
643 for result in results {
647 if let Err(error) = result
648 && error.kind() != std::io::ErrorKind::NotFound
649 {
650 return Err(crate::Error::Io(error));
651 }
652 }
653 Ok(())
654}
655
656#[cfg(test)]
657mod capacity_tests {
658 use super::{MergeCapacity, field_capacity_error};
659
660 #[test]
661 fn merge_capacity_accepts_the_exact_u32_boundary() {
662 let mut capacity = MergeCapacity::default();
663 assert_eq!(capacity.add(u64::from(u32::MAX) - 7), None);
664 assert_eq!(capacity.add(7), None);
665 }
666
667 #[test]
668 fn merge_capacity_rejects_the_first_value_beyond_u32() {
669 let mut capacity = MergeCapacity::default();
670 assert_eq!(capacity.add(u64::from(u32::MAX)), None);
671 assert_eq!(capacity.add(1), Some(u64::from(u32::MAX) + 1));
672 }
673
674 #[test]
675 fn merge_capacity_failure_is_not_source_corruption() {
676 let error = field_capacity_error(
677 7,
678 "body_embedding",
679 "dense vectors",
680 u64::from(u32::MAX) + 1,
681 );
682 assert!(matches!(error, crate::Error::Schema(_)));
683 assert!(error.to_string().contains("lower max_segment_docs"));
684 }
685}