1use super::*;
3use crate::segment::chunk_map::{
4 ChunkMapBuilder, DocLengthsColumn, write_chunk_maps_with_norm_policy,
5};
6use crate::segment::row_map::RowMap;
7use crate::structures::fast_field::{
8 BLOCK_INDEX_ENTRY_SIZE, FastFieldReader, write_fast_field_toc_and_footer,
9};
10use crate::structures::{PositionStreamEncoder, SSTableWriter, TermInfo};
11use std::io::Write;
12
13fn admit(bytes: usize, budget: usize) -> Result<()> {
14 if bytes > budget {
15 return Err(crate::Error::Schema(format!(
16 "compaction scratch requires {bytes} bytes; remaining budget is {budget}"
17 )));
18 }
19 Ok(())
20}
21
22impl SegmentMerger {
23 pub async fn compact<D: DirectoryWriter>(
27 &self,
28 dir: &D,
29 source: &SegmentReader,
30 output: SegmentId,
31 memory_budget: usize,
32 ) -> Result<(SegmentMeta, MergeStats)> {
33 self.ensure_not_cancelled()?;
34 let rows = RowMap::from_visibility(
35 source.num_docs(),
36 source.num_live_docs(),
37 source.alive_docs(),
38 memory_budget / 4,
39 || self.ensure_not_cancelled(),
40 )?;
41 let budget = memory_budget.saturating_sub(rows.memory_bytes());
42 for (field, entry) in self.schema.fields() {
45 if ((entry.indexed && entry.field_type == FieldType::Text)
46 || entry.field_type == FieldType::SparseVector)
47 && !source.row_stats().contains_key(&field.0)
48 {
49 return Err(crate::Error::Schema(format!(
50 "field '{}' lacks lossless row statistics; rebuild this pre-deletion-format segment before compaction",
51 entry.name
52 )));
53 }
54 }
55 let files = SegmentFiles::new(output.0);
56 let mut stats = MergeStats::default();
57 let (field_stats, chunk_maps) = self
58 .compact_text_maps(dir, source, &rows, &files, budget)
59 .await?;
60 let posting_budget =
61 budget.saturating_sub(chunk_maps.values().map(RowMap::memory_bytes).sum::<usize>());
62 stats.terms_processed = self
63 .compact_postings(dir, source, &rows, &chunk_maps, &files, posting_budget)
64 .await?;
65 drop(chunk_maps);
66 stats.fast_bytes = self
67 .compact_columns(dir, &files.fast, source.fast_fields(), &rows, budget)
68 .await?;
69 self.compact_columns(dir, &files.row_stats, source.row_stats(), &rows, budget)
70 .await?;
71 let mut store_out = OffsetWriter::new(dir.streaming_writer_cold(&files.store).await?);
72 let mut store = crate::segment::StoreMerger::new(&mut store_out);
73 store
74 .append_compacted(source.store(), &rows, self.cancellation.as_deref(), budget)
75 .await?;
76 if store.finish()? != rows.len() {
77 return Err(crate::Error::Corruption(
78 "compacted store row count mismatch".into(),
79 ));
80 }
81 stats.store_bytes = store_out.offset() as usize;
82 store_out.finish()?;
83 stats.vectors_bytes = self
84 .compact_vectors(dir, source, &rows, &files, budget)
85 .await?;
86 stats.sparse_bytes = self
87 .compact_sparse(dir, source, &rows, &files, budget)
88 .await?;
89 self.ensure_not_cancelled()?;
90 let meta = SegmentMeta {
91 id: output.0,
92 num_docs: rows.len(),
93 field_stats,
94 };
95 dir.write_durable(&files.meta, &meta.serialize()?).await?;
96 Ok((meta, stats))
97 }
98
99 async fn compact_columns<D: DirectoryWriter>(
100 &self,
101 dir: &D,
102 path: &std::path::Path,
103 columns: &FxHashMap<u32, FastFieldReader>,
104 rows: &RowMap,
105 budget: usize,
106 ) -> Result<usize> {
107 use crate::structures::fast_field::BlockIndexEntry;
108 struct OutputBlock {
109 source: usize,
110 range: std::ops::Range<u32>,
111 header: BlockIndexEntry,
112 copy: bool,
113 cached: Option<Vec<u8>>,
114 }
115 if columns.is_empty() {
116 return Ok(0);
117 }
118 let mut writer = OffsetWriter::new(dir.streaming_writer_cold(path).await?);
119 let mut fields: Vec<_> = columns.keys().copied().collect();
120 fields.sort_unstable();
121 let mut toc = Vec::new();
122 for field in fields {
123 self.ensure_not_cancelled()?;
124 let reader = &columns[&field];
125 if reader.num_docs != rows.physical() {
126 return Err(crate::Error::Corruption(
127 "compaction column row count mismatch".into(),
128 ));
129 }
130 let capacity = reader
131 .blocks()
132 .iter()
133 .map(|block| {
134 let live =
135 rows.count(block.cumulative_docs..block.cumulative_docs + block.num_docs);
136 if live == 0 {
137 0
138 } else if live == block.num_docs {
139 1
140 } else {
141 (live as usize).min((block.num_docs as usize).div_ceil(4096))
142 }
143 })
144 .sum::<usize>();
145 admit(
146 capacity.saturating_mul(std::mem::size_of::<OutputBlock>()),
147 budget / 4,
148 )?;
149 let mut plan = Vec::with_capacity(capacity);
150 let mut cached_bytes = 0usize;
151 for (source, block) in reader.blocks().iter().enumerate() {
152 self.ensure_not_cancelled()?;
153 let live =
154 rows.count(block.cumulative_docs..block.cumulative_docs + block.num_docs);
155 if live == 0 {
156 continue;
157 }
158 if live == block.num_docs {
159 plan.push(OutputBlock {
160 source,
161 range: 0..block.num_docs,
162 copy: true,
163 cached: None,
164 header: BlockIndexEntry {
165 num_docs: live,
166 data_len: block.data.len() as u32,
167 dict_count: block.dict.as_ref().map_or(0, |dict| dict.len()),
168 dict_len: block.raw_dict.len() as u32,
169 },
170 });
171 continue;
172 }
173 for start in (0..block.num_docs).step_by(4096) {
174 self.ensure_not_cancelled()?;
175 let end = start.saturating_add(4096).min(block.num_docs);
176 if rows.count(block.cumulative_docs + start..block.cumulative_docs + end) == 0 {
177 continue;
178 }
179 let bytes = block.compact_range(
180 reader.column_type,
181 reader.multi,
182 start..end,
183 |doc| rows.get(doc).is_some(),
184 budget / 2,
185 )?;
186 let header = BlockIndexEntry::read_from(&mut &bytes[4..])?;
187 let cached = if bytes.capacity() <= (budget / 4).saturating_sub(cached_bytes) {
188 cached_bytes += bytes.capacity();
189 Some(bytes)
190 } else {
191 None
192 };
193 plan.push(OutputBlock {
194 source,
195 range: start..end,
196 header,
197 copy: false,
198 cached,
199 });
200 }
201 }
202 let offset = writer.offset();
203 writer.write_all(&(plan.len() as u32).to_le_bytes())?;
204 for block in &plan {
205 block.header.write_to(&mut writer)?;
206 }
207 for block in plan {
208 self.ensure_not_cancelled()?;
209 let source = &reader.blocks()[block.source];
210 if block.copy {
211 for data in [source.data.as_slice(), source.raw_dict.as_slice()] {
212 for chunk in data.chunks(4 * 1024 * 1024) {
213 self.ensure_not_cancelled()?;
214 writer.write_all(chunk)?;
215 }
216 }
217 } else {
218 let bytes = match block.cached {
219 Some(bytes) => bytes,
220 None => source.compact_range(
221 reader.column_type,
222 reader.multi,
223 block.range,
224 |doc| rows.get(doc).is_some(),
225 budget / 2,
226 )?,
227 };
228 let mut header = Vec::with_capacity(BLOCK_INDEX_ENTRY_SIZE);
229 block.header.write_to(&mut header)?;
230 if bytes.get(4..4 + BLOCK_INDEX_ENTRY_SIZE) != Some(header.as_slice()) {
231 return Err(crate::Error::Corruption(
232 "non-deterministic compacted fast column".into(),
233 ));
234 }
235 writer.write_all(&bytes[4 + BLOCK_INDEX_ENTRY_SIZE..])?;
236 }
237 }
238 toc.push(crate::structures::fast_field::FastFieldTocEntry {
239 field_id: field,
240 column_type: reader.column_type,
241 multi: reader.multi,
242 data_offset: offset,
243 data_len: writer.offset() - offset,
244 num_docs: rows.len(),
245 dict_offset: 0,
246 dict_count: 0,
247 });
248 }
249 let offset = writer.offset();
250 write_fast_field_toc_and_footer(&mut writer, offset, &toc)?;
251 let bytes = writer.offset() as usize;
252 writer.finish()?;
253 Ok(bytes)
254 }
255
256 async fn compact_text_maps<D: DirectoryWriter>(
257 &self,
258 dir: &D,
259 source: &SegmentReader,
260 rows: &RowMap,
261 files: &SegmentFiles,
262 budget: usize,
263 ) -> Result<(FxHashMap<u32, FieldStats>, FxHashMap<u32, RowMap>)> {
264 let mut retained_bytes = 0usize;
265 let mut statistics = FxHashMap::default();
266 let mut chunks = Vec::new();
267 let mut maps = FxHashMap::default();
268 let mut norms = Vec::new();
269 for (field, entry) in self.schema.fields() {
270 if !entry.indexed || entry.field_type != FieldType::Text {
271 continue;
272 }
273 self.ensure_not_cancelled()?;
274 let exact = &source.row_stats()[&field.0];
275 let mut stat = FieldStats::default();
276 for (visited, old) in rows.iter().enumerate() {
277 if visited.is_multiple_of(4096) {
278 self.ensure_not_cancelled()?;
279 }
280 let value = exact.get_u64(old);
281 if value > 0 {
282 stat.doc_count += 1;
283 stat.total_tokens =
284 stat.total_tokens.checked_add(value - 1).ok_or_else(|| {
285 crate::Error::Corruption("compacted token count overflow".into())
286 })?;
287 }
288 }
289 if source.has_text_mapping(field) {
290 if entry.chunked {
291 stat.doc_count = 0;
292 }
293 if let Some(source_map) = source.chunk_map(field) {
294 let map = RowMap::try_new(
295 source_map.num_chunks(),
296 source_map.num_chunks(),
297 |vid| {
298 if vid.is_multiple_of(4096) {
299 self.ensure_not_cancelled()?;
300 }
301 Ok(rows.get(source_map.doc_id(vid)).is_some())
302 },
303 (budget / 2).saturating_sub(retained_bytes),
304 )?;
305 retained_bytes = retained_bytes
306 .saturating_add(map.memory_bytes())
307 .saturating_add(map.len() as usize * 8);
308 admit(retained_bytes, budget / 2)?;
309 let mut builder = ChunkMapBuilder::with_capacity(map.len() as usize);
310 builder.set_document_units(source_map.is_document_map());
311 for (visited, old) in map.iter().enumerate() {
312 if visited.is_multiple_of(4096) {
313 self.ensure_not_cancelled()?;
314 }
315 let (doc, ordinal) = source_map.resolve(old);
316 builder.push(rows.get(doc).unwrap(), ordinal, source_map.length(old))?;
317 }
318 if entry.chunked {
319 stat.doc_count = map.len();
320 }
321 builder.set_total_tokens(stat.total_tokens);
322 chunks.push((field.0, builder));
323 maps.insert(field.0, map);
324 }
325 } else if let Some(lengths) = source.doc_lengths(field) {
326 retained_bytes = retained_bytes.saturating_add(rows.len() as usize * 2);
327 admit(retained_bytes, budget / 2)?;
328 let mut values = Vec::with_capacity(rows.len() as usize);
329 for (visited, old) in rows.iter().enumerate() {
330 if visited.is_multiple_of(4096) {
331 self.ensure_not_cancelled()?;
332 }
333 values.push(lengths.length(old).min(u16::MAX as u32) as u16);
334 }
335 norms.push((field.0, values, stat.total_tokens));
336 }
337 statistics.insert(field.0, stat);
338 }
339 let chunk_refs: Vec<_> = chunks
340 .iter()
341 .filter(|(_, map)| !map.is_empty())
342 .map(|(field, map)| (*field, map))
343 .collect();
344 let norm_refs: Vec<_> = norms
345 .iter()
346 .map(|(field, lengths, total)| DocLengthsColumn {
347 field_id: *field,
348 lengths,
349 total_tokens: *total,
350 })
351 .collect();
352 if !chunk_refs.is_empty() || !norm_refs.is_empty() {
353 let mut writer = dir.streaming_writer_cold(&files.chunks).await?;
354 write_chunk_maps_with_norm_policy(&mut *writer, &chunk_refs, &norm_refs, |field_id| {
355 source
356 .doc_lengths(crate::Field(field_id))
357 .is_some_and(|lengths| lengths.is_quantized())
358 })?;
359 writer.finish()?;
360 }
361 Ok((statistics, maps))
362 }
363
364 async fn compact_postings<D: DirectoryWriter>(
365 &self,
366 dir: &D,
367 source: &SegmentReader,
368 rows: &RowMap,
369 chunks: &FxHashMap<u32, RowMap>,
370 files: &SegmentFiles,
371 budget: usize,
372 ) -> Result<usize> {
373 admit(16 * 1024, budget)?;
376 let budget = budget - 16 * 1024;
377 let mut postings = OffsetWriter::new(dir.streaming_writer_cold(&files.postings).await?);
378 let mut positions = OffsetWriter::new(dir.streaming_writer_cold(&files.positions).await?);
379 let mut terms_out = OffsetWriter::new(dir.streaming_writer_cold(&files.term_dict).await?);
380 let mut terms = SSTableWriter::<_, TermInfo>::with_config(
381 &mut terms_out,
382 crate::structures::SSTableWriterConfig {
383 block_size: self.term_dict_block_size,
384 ..crate::structures::SSTableWriterConfig::from_optimization(self.optimization)
385 },
386 );
387 let mut iter = source.term_dict_iter();
388 let mut count = 0;
389 while let Some((key, info)) = iter.next().await? {
390 self.ensure_not_cancelled()?;
391 let field = crate::Field(u32::from_le_bytes(
392 key.get(..4)
393 .ok_or_else(|| crate::Error::Corruption("invalid term field prefix".into()))?
394 .try_into()
395 .unwrap(),
396 ));
397 let map = chunks.get(&field.0).unwrap_or(rows);
398 let info = if let Some((ids, tfs)) = info.decode_inline() {
399 let entries: Vec<_> = ids
400 .into_iter()
401 .zip(tfs)
402 .filter_map(|(old, tf)| map.get(old).map(|new| (new, tf)))
403 .collect();
404 if entries.is_empty() {
405 continue;
406 }
407 TermInfo::try_inline_iter(entries.len(), entries.into_iter()).ok_or_else(|| {
409 crate::Error::Corruption("compacted inline posting no longer fits".into())
410 })?
411 } else if let Some((offset, len)) = info.external_info() {
412 use crate::structures::postings::{
413 PositionRangeSource, PostingBlockSource, PostingStreamWriter,
414 };
415 let cancellation = self.cancellation.as_deref();
416 let input = PostingBlockSource::open(
417 source.posting_file_range(offset, len)?,
418 budget / 4,
419 cancellation,
420 )
421 .await?;
422 if input.doc_count() != info.doc_freq() || input.doc_count() > map.physical() {
423 return Err(crate::Error::Corruption(
424 "posting source count disagrees with term or row space".into(),
425 ));
426 }
427 let mut source_positions = match info.position_info() {
428 Some((offset, len)) => Some(
429 PositionRangeSource::open(
430 source.position_file_range(offset, len)?,
431 budget / 4,
432 cancellation,
433 )
434 .await?,
435 ),
436 None => None,
437 };
438 let has_positions = source_positions.is_some();
439 if input.has_positions() != has_positions {
440 return Err(crate::Error::Corruption(
441 "posting and position formats disagree".into(),
442 ));
443 }
444 if let Some(positions) = source_positions.as_ref() {
445 let total = if input.len() == 0 {
446 0
447 } else {
448 input.position_span(input.len() - 1)?.end
449 };
450 if positions.total_positions() != total {
451 return Err(crate::Error::Corruption(
452 "posting and position counts disagree".into(),
453 ));
454 }
455 }
456 let off = postings.offset();
457 let pos_off = positions.offset();
458 let mut output = PostingStreamWriter::new(
459 &mut postings,
460 input.len(),
461 has_positions,
462 self.posting_codec,
463 budget / 4,
464 )?;
465 if input.is_compact() && self.posting_codec != crate::structures::PostingCodec::Pfor
466 {
467 output.enable_compact_headers()?;
468 }
469 if input.has_impact_bounds() {
470 output.enable_impact_bounds()?;
471 } else if input.has_ratio_bounds() {
472 output.enable_ratio_bounds()?;
473 }
474 let mut position_output = PositionStreamEncoder::with_budget(
475 &mut positions,
476 budget / 4,
477 self.posting_codec,
478 );
479 if source_positions
480 .as_ref()
481 .is_some_and(PositionRangeSource::is_compact)
482 {
483 position_output = position_output.with_compact_directory();
484 }
485 let mut docs = Vec::with_capacity(crate::structures::postings::POSTING_BLOCK_SIZE);
486 let mut tfs = Vec::with_capacity(crate::structures::postings::POSTING_BLOCK_SIZE);
487 for i in 0..input.len() {
488 self.ensure_not_cancelled()?;
489 let (first, last) = input.bounds(i);
490 if last >= map.physical() {
491 return Err(crate::Error::Corruption(
492 "posting address exceeds compaction map".into(),
493 ));
494 }
495 let live = map.count(first..last + 1);
496 if live == 0 {
497 continue;
498 }
499 let block = input.read_block(i).await?;
500 let span = input.position_span(i)?;
501 if live == last - first + 1 {
502 output.append(&block, map.get(first).expect("live interval"))?;
503 if let Some(positions) = &mut source_positions {
504 positions
505 .append_range(&mut position_output, span, cancellation)
506 .await?;
507 }
508 continue;
509 }
510 if !block.decode_block_into(0, &mut docs, &mut tfs)
511 || docs.first() != Some(&first)
512 || docs.last() != Some(&last)
513 || !docs.windows(2).all(|pair| pair[0] < pair[1])
514 || tfs.contains(&0)
515 {
516 return Err(crate::Error::Corruption(
517 "invalid compacted posting block".into(),
518 ));
519 }
520 let mut cursor = span.start;
521 for (&old, &tf) in docs.iter().zip(&tfs) {
522 if let Some(new) = map.get(old) {
523 let length = source.chunk_map(field).map_or_else(
526 || source.doc_lengths(field).map_or(1, |norm| norm.length(old)),
527 |map| map.length(old),
528 );
529 output.push(new, tf, length)?;
530 if let Some(positions) = &mut source_positions {
531 positions
532 .append_doc(&mut position_output, cursor, tf, cancellation)
533 .await?;
534 }
535 }
536 cursor = cursor.checked_add(u64::from(tf)).ok_or_else(|| {
537 crate::Error::Corruption("position count overflow".into())
538 })?;
539 }
540 if input.has_positions() && cursor != span.end {
541 return Err(crate::Error::Corruption(
542 "position cursor disagrees with term frequencies".into(),
543 ));
544 }
545 }
546 if output.doc_count() == 0 {
547 continue;
548 }
549 if output.doc_count() > input.doc_count() {
550 return Err(crate::Error::Corruption(
551 "compacted posting count exceeds source".into(),
552 ));
553 }
554 let total_positions = output.total_positions();
555 let (docs, len) = output.finish(cancellation)?;
556 if has_positions {
557 self.ensure_not_cancelled()?;
558 let (total, pos_len) = position_output.finish_cancellable(cancellation)?;
559 if total != total_positions {
560 return Err(crate::Error::Corruption(
561 "compacted position count mismatch".into(),
562 ));
563 }
564 TermInfo::external_with_positions(off, len, docs, pos_off, pos_len)
565 } else {
566 TermInfo::external(off, len, docs)
567 }
568 } else {
569 return Err(crate::Error::Corruption(
570 "invalid term posting representation".into(),
571 ));
572 };
573 terms.insert(&key, &info)?;
574 count += 1;
575 }
576 terms.finish()?;
577 terms_out.finish()?;
578 postings.finish()?;
579 if positions.offset() > 0 {
580 positions.finish()?;
581 } else {
582 drop(positions);
583 dir.delete(&files.positions).await?;
584 }
585 Ok(count)
586 }
587}
588
589#[cfg(test)]
590mod tests {
591 use super::*;
592 use crate::structures::BlockPostingList;
593 use crate::structures::fast_field::{FastFieldColumnType, FastFieldWriter};
594
595 #[tokio::test]
596 async fn compaction_preserves_mixed_norm_encodings_per_field() {
597 use crate::segment::builder::{SegmentBuilder, SegmentBuilderConfig};
598 use crate::segment::chunk_map::{
599 read_chunk_maps, write_chunk_maps_with_copied_norms, write_chunk_maps_with_norms,
600 };
601 let mut schema = crate::Schema::builder();
602 let exact = schema.add_text_field("exact", true, false);
603 let quantized = schema.add_text_field("quantized", true, false);
604 let schema = Arc::new(schema.build());
605 let dir = crate::RamDirectory::new();
606 let id = SegmentId::new();
607 let mut builder =
608 SegmentBuilder::new(schema.clone(), SegmentBuilderConfig::default()).unwrap();
609 for _ in 0..16 {
610 let mut doc = crate::Document::new();
611 doc.add_text(exact, "term ".repeat(32));
612 doc.add_text(quantized, "term ".repeat(32));
613 builder.add_document(doc).unwrap();
614 }
615 builder.build(&dir, id, None).await.unwrap();
616 let files = SegmentFiles::new(id.0);
617 let original = read_chunk_maps(
618 dir.open_read(&files.chunks)
619 .await
620 .unwrap()
621 .read_bytes()
622 .await
623 .unwrap(),
624 )
625 .unwrap();
626 let mut bytes = Vec::new();
627 write_chunk_maps_with_norms(
628 &mut bytes,
629 &[],
630 &[DocLengthsColumn {
631 field_id: quantized.0,
632 lengths: &[32; 16],
633 total_tokens: 512,
634 }],
635 true,
636 )
637 .unwrap();
638 let byte_column = read_chunk_maps(crate::directories::OwnedBytes::new(bytes)).unwrap();
639 let mut mixed = Vec::new();
640 write_chunk_maps_with_copied_norms(
641 &mut mixed,
642 &[],
643 &[
644 (exact.0, &original.doc_lengths[&exact.0]),
645 (quantized.0, &byte_column.doc_lengths[&quantized.0]),
646 ],
647 )
648 .unwrap();
649 dir.write(&files.chunks, &mixed).await.unwrap();
652 let source = SegmentReader::open(&dir, id, schema.clone(), 4)
653 .await
654 .unwrap();
655 let rows = RowMap::new(16, 16, |doc| doc % 2 == 0, 1024).unwrap();
656 let output = SegmentFiles::new(SegmentId::new().0);
657 SegmentMerger::new(schema)
658 .compact_text_maps(&dir, &source, &rows, &output, 1024 * 1024)
659 .await
660 .unwrap();
661 let columns = read_chunk_maps(
662 dir.open_read(&output.chunks)
663 .await
664 .unwrap()
665 .read_bytes()
666 .await
667 .unwrap(),
668 )
669 .unwrap();
670 for (field, encoded) in [(exact, false), (quantized, true)] {
671 let lengths = &columns.doc_lengths[&field.0];
672 assert_eq!(lengths.is_quantized(), encoded);
673 assert_eq!(lengths.num_docs(), 8);
674 assert_eq!(lengths.total_tokens(), 256);
675 assert!((0..8).all(|doc| lengths.length(doc) == 32));
676 }
677 }
678
679 #[tokio::test]
680 async fn frequent_terms_compact_with_block_sized_decoding_and_preserve_positions() {
681 check_frequent_compaction(false).await;
682 }
683
684 #[tokio::test]
685 async fn compact_layout_survives_dense_and_sparse_row_compaction() {
686 check_frequent_compaction(true).await;
687 }
688
689 async fn check_frequent_compaction(compact: bool) {
690 use crate::directories::{Directory, RamDirectory};
691 use crate::structures::{TERMINATED, TermPositions};
692 let mut schema = crate::SchemaBuilder::default();
693 let body = schema.add_text_field("body", true, false);
694 schema.set_positions(body, crate::dsl::PositionMode::TokenPosition);
695 let schema = schema.build();
696 let dir = RamDirectory::new();
697 let index = crate::Index::create(
698 dir.clone(),
699 schema.clone(),
700 crate::IndexConfig {
701 compact_text: compact,
702 quantized_norms: compact,
703 merge_policy: Box::new(crate::NoMergePolicy),
704 ..Default::default()
705 },
706 )
707 .await
708 .unwrap();
709 let mut writer = index.writer();
710 for _ in 0..8192 {
711 let mut doc = crate::Document::new();
712 doc.add_text(body, "common common common common");
713 writer.add_document(doc).unwrap();
714 }
715 writer.commit().await.unwrap();
716 let reader = index.reader().await.unwrap();
717 let searcher = reader.searcher().await.unwrap();
718 for sparse in [false, true] {
719 let rows = RowMap::new(
720 8192,
721 8192,
722 |doc| {
723 if sparse {
724 doc % 129 == 0
725 } else {
726 doc >= 256 && doc % 257 != 0
727 }
728 },
729 4096,
730 )
731 .unwrap();
732 let merger = SegmentMerger::new(Arc::new(schema.clone()));
733 let files = SegmentFiles::new(SegmentId::new().0);
734 merger
736 .compact_postings(
737 &dir,
738 &searcher.segment_readers()[0],
739 &rows,
740 &FxHashMap::default(),
741 &files,
742 64 * 1024,
743 )
744 .await
745 .unwrap();
746 let posting_bytes = dir
747 .open_read(&files.postings)
748 .await
749 .unwrap()
750 .read_bytes()
751 .await
752 .unwrap();
753 if compact {
754 let flags = u32::from_le_bytes(
755 posting_bytes[posting_bytes.len() - 12..posting_bytes.len() - 8]
756 .try_into()
757 .unwrap(),
758 );
759 assert_ne!(flags & 64, 0, "compaction retains separated headers");
760 }
761 let postings = BlockPostingList::deserialize(posting_bytes.as_slice()).unwrap();
762 assert_eq!(postings.doc_count(), rows.len());
763 if sparse {
764 assert_eq!(
765 postings.num_blocks(),
766 1,
767 "sparse survivors should share a rebuilt block"
768 );
769 }
770 let position_bytes = dir
771 .open_read(&files.positions)
772 .await
773 .unwrap()
774 .read_bytes()
775 .await
776 .unwrap();
777 if compact {
778 assert_eq!(&position_bytes[position_bytes.len() - 4..], b"POS6");
779 }
780 let positions = TermPositions::open(position_bytes).unwrap();
781 let mut cursor = postings.iterator();
782 let mut scratch = Vec::new();
783 let mut values = Vec::new();
784 for doc in 0..rows.len() {
785 assert_eq!(cursor.doc(), doc);
786 assert_eq!(cursor.term_freq(), 4);
787 assert!(positions.positions_into(
788 cursor.position_cursor(),
789 4,
790 &mut scratch,
791 &mut values
792 ));
793 assert_eq!(values, [0, 1, 2, 3]);
794 cursor.advance();
795 }
796 assert_eq!(cursor.doc(), TERMINATED);
797 }
798 }
799
800 #[tokio::test]
801 async fn intact_fast_blocks_keep_encoded_bytes_when_earlier_rows_are_deleted() {
802 let n = 8193;
803 let mut column = FastFieldWriter::new_numeric(FastFieldColumnType::U64);
804 for i in 0..n {
805 column.add_u64(i, u64::from(i) * 17);
806 }
807 let mut encoded = Vec::new();
808 let (mut toc, _) = column.serialize(&mut encoded, 0).unwrap();
809 let header = &encoded[4..4 + BLOCK_INDEX_ENTRY_SIZE];
810 let payload = &encoded[4 + BLOCK_INDEX_ENTRY_SIZE..];
811 let mut stacked = 2u32.to_le_bytes().to_vec();
812 stacked.extend_from_slice(header);
813 stacked.extend_from_slice(header);
814 stacked.extend_from_slice(payload);
815 stacked.extend_from_slice(payload);
816 toc.num_docs = 2 * n;
817 toc.data_len = stacked.len() as u64;
818 let source =
819 FastFieldReader::open(&crate::directories::OwnedBytes::new(stacked), &toc).unwrap();
820 let columns = FxHashMap::from_iter([(0, source)]);
821 let rows = RowMap::new(2 * n, n, |doc| doc >= n, 1024 * 1024).unwrap();
822 let dir = crate::directories::RamDirectory::new();
823 let merger = SegmentMerger::new(Arc::new(crate::SchemaBuilder::default().build()));
824 let path = std::path::Path::new("copied.fast");
825 merger
826 .compact_columns(&dir, path, &columns, &rows, 4 * 1024 * 1024)
827 .await
828 .unwrap();
829 let bytes = dir
830 .open_read(path)
831 .await
832 .unwrap()
833 .read_bytes()
834 .await
835 .unwrap();
836 let (offset, count) =
837 crate::structures::fast_field::read_fast_field_footer(bytes.as_slice()).unwrap();
838 let toc =
839 crate::structures::fast_field::read_fast_field_toc(bytes.as_slice(), offset, count)
840 .unwrap();
841 assert_eq!(toc[0].data_len as usize, encoded.len());
842 assert_eq!(&bytes.as_slice()[..encoded.len()], encoded.as_slice());
843 }
844
845 #[tokio::test]
846 async fn cached_and_reencoded_column_chunks_have_identical_bytes_and_missing_rows() {
847 let count = 2 * (32 * 4096 + 1);
850 let mut column = FastFieldWriter::new_numeric(FastFieldColumnType::U64);
851 for doc in 0..count {
852 if doc % 97 != 0 {
853 column.add_u64(
854 doc,
855 u64::from(doc)
856 .wrapping_mul(0x9e3779b97f4a7c15)
857 .rotate_left(23),
858 );
859 }
860 }
861 column.pad_to(count);
862 let mut encoded = Vec::new();
863 let (toc, _) = column.serialize(&mut encoded, 0).unwrap();
864 let source =
865 FastFieldReader::open(&crate::directories::OwnedBytes::new(encoded), &toc).unwrap();
866 let columns = FxHashMap::from_iter([(0, source)]);
867 let rows = RowMap::new(count, count / 2, |doc| doc % 2 == 1, 8 * 1024 * 1024).unwrap();
868 let dir = crate::directories::RamDirectory::new();
869 let merger = SegmentMerger::new(Arc::new(crate::SchemaBuilder::default().build()));
870 let mut outputs = Vec::new();
871 for (name, budget) in [
872 ("partial.fast", 4 * 1024 * 1024),
873 ("cached.fast", 16 * 1024 * 1024),
874 ] {
875 let path = std::path::Path::new(name);
876 merger
877 .compact_columns(&dir, path, &columns, &rows, budget)
878 .await
879 .unwrap();
880 outputs.push(
881 dir.open_read(path)
882 .await
883 .unwrap()
884 .read_bytes()
885 .await
886 .unwrap(),
887 );
888 }
889 assert!(
890 outputs[0].len() > 1024 * 1024,
891 "fixture no longer exercises cache overflow"
892 );
893 assert_eq!(outputs[0].as_slice(), outputs[1].as_slice());
894 let (offset, count) =
895 crate::structures::fast_field::read_fast_field_footer(outputs[0].as_slice()).unwrap();
896 let toc = crate::structures::fast_field::read_fast_field_toc(
897 outputs[0].as_slice(),
898 offset,
899 count,
900 )
901 .unwrap();
902 let actual = FastFieldReader::open(&outputs[0], &toc[0]).unwrap();
903 assert_eq!(actual.num_docs, rows.len());
904 for (new, old) in rows.iter().enumerate() {
905 assert_eq!(actual.get_u64(new as u32), columns[&0].get_u64(old));
906 assert_eq!(actual.has_value(new as u32), columns[&0].has_value(old));
907 }
908 }
909}