1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
//! Document store merge.
//!
//! Merges store blocks from multiple segments, handling both raw block
//! copying (when no dictionary compression) and recompression paths.
use super::OffsetWriter;
use super::SegmentMerger;
use crate::Result;
use crate::segment::reader::SegmentReader;
use crate::segment::store::StoreMerger;
impl SegmentMerger {
/// Merge document stores from all source segments into a single output.
///
/// Uses raw block copying when possible (no dictionary compression),
/// falling back to recompression when segments use dictionaries.
/// Returns (store_num_docs) so the caller can verify against metadata.
pub(super) async fn merge_store(
&self,
segments: &[SegmentReader],
store_writer: &mut OffsetWriter,
) -> Result<u32> {
// Pre-check: verify each source segment's store num_docs matches its metadata.
// A mismatch means the source segment is corrupt (e.g. store blocks were lost
// during initial build). Proceeding would desynchronize postings and store
// doc_id spaces in the merged output, causing progressive document loss.
for segment in segments {
let meta_docs = segment.num_docs();
let store_docs = segment.store().num_docs();
if meta_docs != store_docs {
log::error!(
"[merge_store] index={} SOURCE MISMATCH segment {:016x}: meta.num_docs={}, store.num_docs={}",
self.schema.index_label(),
segment.meta().id,
meta_docs,
store_docs
);
return Err(crate::Error::Io(std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!(
"Source segment {:016x} has store/meta mismatch: store={}, meta={}",
segment.meta().id,
store_docs,
meta_docs
),
)));
}
}
let mut store_merger = StoreMerger::new(store_writer);
for segment in segments {
self.ensure_not_cancelled()?;
if segment.store_has_dict() {
let result = store_merger
.append_store_recompressing(segment.store(), self.cancellation.as_deref())
.await;
self.ensure_not_cancelled()?;
result.map_err(crate::Error::Io)?;
} else {
let raw_blocks = segment.store_raw_blocks();
let data_slice = segment.store_data_slice();
let result = store_merger
.append_store(data_slice, &raw_blocks, self.cancellation.as_deref())
.await;
self.ensure_not_cancelled()?;
result.map_err(crate::Error::Io)?;
}
}
let store_num_docs = store_merger.finish()?;
Ok(store_num_docs)
}
}