1mod 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
24fn 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#[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#[derive(Debug, Clone, Default)]
73pub struct MergeStats {
74 pub terms_processed: usize,
76 pub term_dict_bytes: usize,
78 pub postings_bytes: usize,
80 pub store_bytes: usize,
82 pub vectors_bytes: usize,
84 pub sparse_bytes: usize,
86 pub bp_converged: bool,
92 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
112pub use super::types::TrainedVectorStructures;
114
115pub(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
129pub(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 log::warn!(
162 "[merge] failed to remove temporary sparse section {:?}: {}",
163 path,
164 error,
165 );
166 }
167 Ok(())
168}
169
170pub struct SegmentMerger {
172 schema: Arc<Schema>,
173 reorder_bmp: bool,
177 background_pool: Option<Arc<rayon::ThreadPool>>,
181 granularity: crate::segment::reorder::BpGranularity,
185 bp_budget: crate::segment::BpBudget,
190 bp_memory_budget: usize,
192 reorder_permits: Option<Arc<ReorderConcurrencyGate>>,
195 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 pub fn with_bmp_reorder(mut self, reorder: bool) -> Self {
216 self.reorder_bmp = reorder;
217 self
218 }
219
220 pub fn with_background_pool(mut self, pool: Option<Arc<rayon::ThreadPool>>) -> Self {
222 self.background_pool = pool;
223 self
224 }
225
226 pub fn with_granularity(mut self, granularity: crate::segment::reorder::BpGranularity) -> Self {
228 self.granularity = granularity;
229 self
230 }
231
232 pub fn with_bp_budget(mut self, budget: crate::segment::BpBudget) -> Self {
234 self.bp_budget = budget;
235 self
236 }
237
238 pub fn with_bp_memory_budget(mut self, bytes: usize) -> Self {
240 self.bp_memory_budget = bytes;
241 self
242 }
243
244 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 fn validate_merge_capacities(&self, segments: &[SegmentReader]) -> Result<()> {
259 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 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 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 let merge_start = std::time::Instant::now();
434
435 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 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 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 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 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 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
608pub 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 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}