Skip to main content

risingwave_storage/hummock/sstable/
builder.rs

1// Copyright 2022 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::{BTreeMap, BTreeSet, HashMap};
16use std::mem;
17use std::sync::Arc;
18use std::time::SystemTime;
19
20use await_tree::{InstrumentAwait, SpanExt};
21use bytes::{Bytes, BytesMut};
22use risingwave_common::catalog::TableId;
23use risingwave_common::hash::VirtualNode;
24use risingwave_common::util::row_serde::OrderedRowSerde;
25use risingwave_hummock_sdk::key::{FullKey, MAX_KEY_LEN, TABLE_PREFIX_LEN, UserKey, user_key};
26use risingwave_hummock_sdk::key_range::KeyRange;
27use risingwave_hummock_sdk::sstable_info::{SstableInfo, SstableInfoInner, VnodeStatistics};
28use risingwave_hummock_sdk::table_stats::{TableStats, TableStatsMap};
29use risingwave_hummock_sdk::{HummockEpoch, HummockSstableObjectId, LocalSstableInfo};
30use risingwave_pb::hummock::{PbSstableFilterLayout, PbSstableFilterType};
31
32use super::utils::CompressionAlgorithm;
33use super::{
34    BlockBuilder, BlockBuilderOptions, BlockMeta, DEFAULT_BLOCK_SIZE, DEFAULT_ENTRY_SIZE,
35    DEFAULT_RESTART_INTERVAL, SstableMeta, SstableWriter, VERSION,
36};
37use crate::compaction_catalog_manager::{
38    CompactionCatalogAgent, CompactionCatalogAgentRef, FilterKeyExtractorImpl,
39    FullKeyFilterKeyExtractor,
40};
41use crate::hummock::sstable::{
42    DEFAULT_FILTER_HASH_PREALLOC_KEY_COUNT_CAP, FilterBuilder, FilterBuilderOptions, utils,
43};
44use crate::hummock::value::HummockValue;
45use crate::hummock::{
46    Block, BlockHolder, BlockIterator, HummockResult, MemoryLimiter, Xor16FilterBuilder,
47    try_shorten_block_smallest_key,
48};
49use crate::monitor::CompactorMetrics;
50use crate::opts::StorageOpts;
51
52pub const DEFAULT_SSTABLE_SIZE: usize = 4 * 1024 * 1024;
53pub const DEFAULT_BLOOM_FALSE_POSITIVE: f64 = 0.001;
54pub const DEFAULT_MAX_SST_SIZE: u64 = 512 * 1024 * 1024;
55pub const MIN_BLOCK_SIZE: usize = 8 * 1024;
56
57#[derive(Clone, Debug)]
58pub struct SstableBuilderOptions {
59    /// Approximate sstable capacity.
60    pub capacity: usize,
61    /// Approximate block capacity.
62    pub block_capacity: usize,
63    /// Restart point interval.
64    pub restart_interval: usize,
65    /// Deprecated and ignored by SST filter builders; kept for backward compatibility.
66    pub bloom_false_positive: f64,
67    /// Compression algorithm.
68    pub compression_algorithm: CompressionAlgorithm,
69    pub max_sst_size: u64,
70    /// If set, block metadata keys will be shortened when their length exceeds this threshold.
71    pub shorten_block_meta_key_threshold: Option<usize>,
72    /// Max bytes for vnode key-range hints in SST metadata. None disables collection.
73    pub max_vnode_key_range_bytes: Option<usize>,
74    /// Estimated key count for one output SST. Used only as a filter-builder capacity hint.
75    pub estimated_output_key_count: Option<usize>,
76    /// Upper bound for the initial key-hash buffer allocation in plain filter builders.
77    pub filter_hash_prealloc_key_count_cap: usize,
78}
79
80impl From<&StorageOpts> for SstableBuilderOptions {
81    fn from(options: &StorageOpts) -> SstableBuilderOptions {
82        let capacity: usize = (options.sstable_size_mb as usize) * (1 << 20);
83        SstableBuilderOptions {
84            capacity,
85            block_capacity: (options.block_size_kb as usize) * (1 << 10),
86            restart_interval: DEFAULT_RESTART_INTERVAL,
87            bloom_false_positive: options.bloom_false_positive,
88            compression_algorithm: CompressionAlgorithm::None,
89            max_sst_size: options.compactor_max_sst_size,
90            shorten_block_meta_key_threshold: options.shorten_block_meta_key_threshold,
91            max_vnode_key_range_bytes: None,
92            estimated_output_key_count: None,
93            filter_hash_prealloc_key_count_cap: DEFAULT_FILTER_HASH_PREALLOC_KEY_COUNT_CAP,
94        }
95    }
96}
97
98impl Default for SstableBuilderOptions {
99    fn default() -> Self {
100        Self {
101            capacity: DEFAULT_SSTABLE_SIZE,
102            block_capacity: DEFAULT_BLOCK_SIZE,
103            restart_interval: DEFAULT_RESTART_INTERVAL,
104            bloom_false_positive: DEFAULT_BLOOM_FALSE_POSITIVE,
105            compression_algorithm: CompressionAlgorithm::None,
106            max_sst_size: DEFAULT_MAX_SST_SIZE,
107            shorten_block_meta_key_threshold: None,
108            max_vnode_key_range_bytes: None,
109            estimated_output_key_count: None,
110            filter_hash_prealloc_key_count_cap: DEFAULT_FILTER_HASH_PREALLOC_KEY_COUNT_CAP,
111        }
112    }
113}
114
115impl SstableBuilderOptions {
116    pub fn estimated_output_key_count(&self) -> usize {
117        self.estimated_output_key_count
118            .unwrap_or(self.capacity / DEFAULT_ENTRY_SIZE + 1)
119    }
120
121    pub fn filter_builder_options(&self) -> FilterBuilderOptions {
122        FilterBuilderOptions {
123            estimated_key_count: self.estimated_output_key_count(),
124            estimated_block_count: self.capacity / self.block_capacity + 1,
125            hash_prealloc_key_count_cap: self.filter_hash_prealloc_key_count_cap,
126        }
127    }
128}
129
130pub struct SstableBuilderOutput<WO> {
131    pub sst_info: LocalSstableInfo,
132    pub writer_output: WO,
133    pub stats: SstableBuilderOutputStats,
134}
135
136pub struct SstableBuilder<W: SstableWriter, F: FilterBuilder> {
137    /// Options.
138    options: SstableBuilderOptions,
139    /// Data writer.
140    writer: W,
141    /// Current block builder.
142    block_builder: BlockBuilder,
143
144    compaction_catalog_agent_ref: CompactionCatalogAgentRef,
145    /// Block metadata vec.
146    block_metas: Vec<BlockMeta>,
147
148    /// `table_id` of added keys.
149    table_ids: BTreeSet<TableId>,
150    last_full_key: Vec<u8>,
151    /// Encoded key buffer, swapped with `last_full_key` after each entry.
152    raw_key: Vec<u8>,
153    /// Reusable buffer for the encoded value.
154    raw_value: BytesMut,
155    last_table_id: Option<TableId>,
156    sst_object_id: HummockSstableObjectId,
157
158    /// Per table stats.
159    table_stats: TableStatsMap,
160    /// `last_table_stats` accumulates stats for `last_table_id` and finalizes it in `table_stats`
161    /// by `finalize_last_table_stats`
162    last_table_stats: TableStats,
163
164    filter_builder: F,
165
166    epoch_set: BTreeSet<u64>,
167    memory_limiter: Option<Arc<MemoryLimiter>>,
168
169    block_size_vec: Vec<usize>, // for statistics
170    vnode_range_collector: Option<VnodeUserKeyRangeCollector>,
171}
172
173impl<W: SstableWriter> SstableBuilder<W, Xor16FilterBuilder> {
174    pub fn for_test(
175        sstable_id: u64,
176        writer: W,
177        options: SstableBuilderOptions,
178        table_id_to_vnode: HashMap<impl Into<TableId>, usize>,
179        table_id_to_watermark_serde: HashMap<
180            impl Into<TableId>,
181            Option<(OrderedRowSerde, OrderedRowSerde, usize)>,
182        >,
183    ) -> Self {
184        let compaction_catalog_agent_ref = Arc::new(CompactionCatalogAgent::new(
185            FilterKeyExtractorImpl::FullKey(FullKeyFilterKeyExtractor),
186            table_id_to_vnode
187                .into_iter()
188                .map(|(table_id, v)| (table_id.into(), v))
189                .collect(),
190            table_id_to_watermark_serde
191                .into_iter()
192                .map(|(table_id, v)| (table_id.into(), v))
193                .collect(),
194            HashMap::default(),
195        ));
196
197        Self::new(
198            sstable_id,
199            writer,
200            Xor16FilterBuilder::create(options.filter_builder_options()),
201            options,
202            compaction_catalog_agent_ref,
203            None,
204        )
205    }
206}
207
208impl<W: SstableWriter, F: FilterBuilder> SstableBuilder<W, F> {
209    pub fn new(
210        sst_object_id: impl Into<HummockSstableObjectId>,
211        writer: W,
212        filter_builder: F,
213        options: SstableBuilderOptions,
214        compaction_catalog_agent_ref: CompactionCatalogAgentRef,
215        memory_limiter: Option<Arc<MemoryLimiter>>,
216    ) -> Self {
217        let sst_object_id = sst_object_id.into();
218        Self {
219            vnode_range_collector: VnodeUserKeyRangeCollector::with_limit(
220                options.max_vnode_key_range_bytes,
221            ),
222            options: options.clone(),
223            writer,
224            block_builder: BlockBuilder::new(BlockBuilderOptions {
225                capacity: options.block_capacity,
226                restart_interval: options.restart_interval,
227                compression_algorithm: options.compression_algorithm,
228            }),
229            filter_builder,
230            block_metas: Vec::with_capacity(options.capacity / options.block_capacity + 1),
231            table_ids: BTreeSet::new(),
232            last_table_id: None,
233            raw_key: Vec::new(),
234            raw_value: BytesMut::new(),
235            last_full_key: vec![],
236            sst_object_id,
237            compaction_catalog_agent_ref,
238            table_stats: Default::default(),
239            last_table_stats: Default::default(),
240            epoch_set: BTreeSet::default(),
241            memory_limiter,
242            block_size_vec: Vec::new(),
243        }
244    }
245
246    /// Add kv pair to sstable.
247    pub async fn add_for_test(
248        &mut self,
249        full_key: FullKey<&[u8]>,
250        value: HummockValue<&[u8]>,
251    ) -> HummockResult<()> {
252        self.add(full_key, value).await
253    }
254
255    pub fn current_block_size(&self) -> usize {
256        self.block_builder.approximate_len()
257    }
258
259    /// Append a block, compressing uncompressed input with the output codec.
260    /// Returns its written length, or `None` if merged into a small pending block.
261    pub async fn add_raw_block(
262        &mut self,
263        mut buf: Bytes,
264        filter_data: Vec<u8>,
265        smallest_key: FullKey<Vec<u8>>,
266        largest_key: Vec<u8>,
267        mut meta: BlockMeta,
268    ) -> HummockResult<Option<usize>> {
269        let table_id = smallest_key.user_key.table_id;
270        if self.last_table_id.is_none() || self.last_table_id.unwrap() != table_id {
271            if !self.block_builder.is_empty() {
272                // Try to finish the previous `Block`` when the `table_id` is switched, making sure that the data in the `Block` doesn't span two `table_ids`.
273                self.build_block()
274                    .instrument_await("sstable_build_block_on_raw_table_switch".verbose())
275                    .await?;
276            }
277
278            self.table_ids.insert(table_id);
279            self.finalize_last_table_stats();
280            self.last_table_id = Some(table_id);
281        }
282
283        if !self.block_builder.is_empty() {
284            let min_block_size = std::cmp::min(MIN_BLOCK_SIZE, self.options.block_capacity / 4);
285
286            // If the previous block is too small, we should merge it into the previous block.
287            if self.block_builder.approximate_len() < min_block_size {
288                let block = Block::decode(buf, meta.uncompressed_size as usize)?;
289                let mut iter = BlockIterator::new(BlockHolder::from_owned_block(Box::new(block)));
290                iter.seek_to_first();
291                while iter.is_valid() {
292                    let value = HummockValue::from_slice(iter.value()).unwrap_or_else(|_| {
293                        panic!(
294                            "decode failed for fast compact sst_id {} block_idx {} last_table_id {:?}",
295                            self.sst_object_id, self.block_metas.len(), self.last_table_id
296                        )
297                    });
298                    self.add_impl(iter.key(), value, false)
299                        .instrument_await("sstable_add_raw_block_fallback_add".verbose())
300                        .await?;
301                    iter.next();
302                }
303                return Ok(None);
304            }
305
306            self.build_block()
307                .instrument_await("sstable_build_block_before_raw_block".verbose())
308                .await?;
309        }
310        // Decide whether to coalesce first, avoiding compression followed by immediate
311        // decompression. Already compressed blocks retain their original codec.
312        let algorithm = Block::get_algorithm(&buf)?;
313        if algorithm == CompressionAlgorithm::None
314            && algorithm != self.options.compression_algorithm
315        {
316            buf = BlockBuilder::compress_block(buf, self.options.compression_algorithm)?;
317            meta.len = buf.len() as u32;
318        }
319        self.last_full_key = largest_key;
320        assert_eq!(
321            meta.len as usize,
322            buf.len(),
323            "meta {} buf {} last_table_id {:?}",
324            meta.len,
325            buf.len(),
326            self.last_table_id
327        );
328        meta.offset = self.writer.data_len() as u32;
329        let block_len = buf.len();
330        self.block_metas.push(meta);
331        self.filter_builder.add_raw_data(filter_data);
332        let block_meta = self.block_metas.last_mut().unwrap();
333        self.writer
334            .write_block_bytes(buf, block_meta)
335            .instrument_await("sstable_write_raw_block_bytes".verbose())
336            .await?;
337
338        Ok(Some(block_len))
339    }
340
341    /// Add kv pair to sstable.
342    pub async fn add(
343        &mut self,
344        full_key: FullKey<&[u8]>,
345        value: HummockValue<&[u8]>,
346    ) -> HummockResult<()> {
347        self.add_impl(full_key, value, true).await
348    }
349
350    /// Add kv pair to sstable.
351    async fn add_impl(
352        &mut self,
353        full_key: FullKey<&[u8]>,
354        value: HummockValue<&[u8]>,
355        could_switch_block: bool,
356    ) -> HummockResult<()> {
357        const LARGE_KEY_LEN: usize = MAX_KEY_LEN >> 1;
358
359        let table_key_len = full_key.user_key.table_key.as_ref().len();
360        let table_value_len = match &value {
361            HummockValue::Put(t) => t.len(),
362            HummockValue::Delete => 0,
363        };
364        let large_value_len = self.options.max_sst_size as usize / 10;
365        let large_key_value_len = self.options.max_sst_size as usize / 2;
366        if table_key_len >= LARGE_KEY_LEN
367            || table_value_len > large_value_len
368            || table_key_len + table_value_len > large_key_value_len
369        {
370            let table_id = full_key.user_key.table_id;
371            tracing::warn!(
372                "A large key/value (table_id={}, key len={}, value len={}, epoch={}, spill offset={}) is added to block",
373                table_id,
374                table_key_len,
375                table_value_len,
376                full_key.epoch_with_gap.pure_epoch(),
377                full_key.epoch_with_gap.offset(),
378            );
379        }
380
381        // TODO: refine me
382        full_key.encode_into(&mut self.raw_key);
383        value.encode(&mut self.raw_value);
384        let is_new_user_key = self.last_full_key.is_empty()
385            || !user_key(&self.raw_key).eq(user_key(self.last_full_key.as_slice()));
386        let table_id = full_key.user_key.table_id;
387        let is_new_table = self.last_table_id != Some(table_id);
388        let current_block_size = self.current_block_size();
389        let is_block_full = current_block_size >= self.options.block_capacity
390            || (current_block_size > self.options.block_capacity / 4 * 3
391                && current_block_size + self.raw_value.len() + self.raw_key.len()
392                    > self.options.block_capacity);
393
394        if is_new_table {
395            assert!(
396                could_switch_block,
397                "is_new_user_key {} sst_id {} block_idx {} table_id {} last_table_id {:?} full_key {:?}",
398                is_new_user_key,
399                self.sst_object_id,
400                self.block_metas.len(),
401                table_id,
402                self.last_table_id,
403                full_key
404            );
405            self.table_ids.insert(table_id);
406            self.finalize_last_table_stats();
407            self.last_table_id = Some(table_id);
408            if !self.block_builder.is_empty() {
409                self.build_block()
410                    .instrument_await("sstable_build_block_on_table_switch".verbose())
411                    .await?;
412            }
413        } else if is_block_full && could_switch_block {
414            self.build_block()
415                .instrument_await("sstable_build_block_on_block_full".verbose())
416                .await?;
417        }
418        self.last_table_stats.total_key_count += 1;
419        self.epoch_set.insert(full_key.epoch_with_gap.pure_epoch());
420
421        // Rotate block builder if the previous one has been built.
422        if self.block_builder.is_empty() {
423            let smallest_key = if let Some(threshold) =
424                self.options.shorten_block_meta_key_threshold
425                && !self.last_full_key.is_empty()
426                && full_key.encoded_len() >= threshold
427            {
428                let prev = FullKey::decode(&self.last_full_key);
429                if let Some(shortened) = try_shorten_block_smallest_key(&prev, &full_key) {
430                    shortened.encode()
431                } else {
432                    full_key.encode()
433                }
434            } else {
435                full_key.encode()
436            };
437
438            self.block_metas.push(BlockMeta {
439                offset: utils::checked_into_u32(self.writer.data_len()).unwrap_or_else(|_| {
440                    panic!(
441                        "WARN overflow can't convert writer_data_len {} into u32 sst_id {} block_idx {} tables {:?}",
442                        self.writer.data_len(),
443                        self.sst_object_id,
444                        self.block_metas.len(),
445                        self.table_ids,
446                    )
447                }),
448                len: 0,
449                smallest_key,
450                uncompressed_size: 0,
451                total_key_count: 0,
452                stale_key_count: 0,
453            });
454        }
455
456        let filter_key = self
457            .compaction_catalog_agent_ref
458            .extract(user_key(&self.raw_key));
459        // Add SST filter check.
460        if !filter_key.is_empty() {
461            self.filter_builder
462                .add_key(filter_key, table_id.as_raw_id());
463        }
464        // Use pre-encoded key to avoid redundant encoding
465        self.block_builder.add(
466            table_id,
467            &self.raw_key[TABLE_PREFIX_LEN..],
468            self.raw_value.as_ref(),
469        );
470        self.block_metas.last_mut().unwrap().total_key_count += 1;
471        if !is_new_user_key || value.is_delete() {
472            self.block_metas.last_mut().unwrap().stale_key_count += 1;
473        }
474        self.last_table_stats.total_key_size += full_key.encoded_len() as i64;
475        self.last_table_stats.total_value_size += value.encoded_len() as i64;
476
477        if let Some(collector) = self.vnode_range_collector.as_mut() {
478            collector.observe_key(
479                VirtualNode::from_index(full_key.user_key.get_vnode_id()),
480                &self.raw_key,
481                &self.last_full_key,
482            );
483        }
484
485        mem::swap(&mut self.last_full_key, &mut self.raw_key);
486
487        self.raw_key.clear();
488        self.raw_value.clear();
489        Ok(())
490    }
491
492    /// Finish building sst.
493    ///
494    /// # Format
495    ///
496    /// data:
497    ///
498    /// ```plain
499    /// | Block 0 | ... | Block N-1 | N (4B) |
500    /// ```
501    pub async fn finish(mut self) -> HummockResult<SstableBuilderOutput<W::Output>> {
502        let smallest_key = if self.block_metas.is_empty() {
503            vec![]
504        } else {
505            self.block_metas[0].smallest_key.clone()
506        };
507        let largest_key = self.last_full_key.clone();
508        self.finalize_last_table_stats();
509
510        // Vnode key-range hints are only supported for single-table SSTs.
511        // Multi-table SST scenarios should not enable max_vnode_key_range_bytes in config.
512        assert!(
513            self.table_ids.len() <= 1 || self.vnode_range_collector.is_none(),
514            "vnode key-range hints are only supported for single-table SSTs, found {} tables",
515            self.table_ids.len()
516        );
517
518        self.build_block()
519            .instrument_await("sstable_finish_build_block".verbose())
520            .await?;
521        let right_exclusive = false;
522        let meta_offset = self.writer.data_len() as u64;
523
524        let filter_data = self.filter_builder.finish(self.memory_limiter.clone());
525        let (filter_type, filter_layout) = if filter_data.is_none() {
526            (
527                PbSstableFilterType::SstableFilterNone,
528                PbSstableFilterLayout::Unspecified,
529            )
530        } else if self.filter_builder.support_blocked_raw_data() {
531            (
532                self.filter_builder.filter_type(),
533                PbSstableFilterLayout::Blocked,
534            )
535        } else {
536            (
537                self.filter_builder.filter_type(),
538                PbSstableFilterLayout::Plain,
539            )
540        };
541
542        let total_key_count = self
543            .block_metas
544            .iter()
545            .map(|block_meta| block_meta.total_key_count as u64)
546            .sum::<u64>();
547        let stale_key_count = self
548            .block_metas
549            .iter()
550            .map(|block_meta| block_meta.stale_key_count as u64)
551            .sum::<u64>();
552        let uncompressed_file_size = self
553            .block_metas
554            .iter()
555            .map(|block_meta| block_meta.uncompressed_size as u64)
556            .sum::<u64>();
557
558        #[expect(deprecated)]
559        let mut meta = SstableMeta {
560            block_metas: self.block_metas,
561            bloom_filter: filter_data.unwrap_or_default(),
562            estimated_size: 0,
563            key_count: utils::checked_into_u32(total_key_count).unwrap_or_else(|_| {
564                panic!(
565                    "WARN overflow can't convert total_key_count {} into u32 tables {:?}",
566                    total_key_count, self.table_ids,
567                )
568            }),
569            smallest_key,
570            largest_key,
571            version: VERSION,
572            meta_offset,
573            monotonic_tombstone_events: vec![],
574        };
575
576        let meta_encode_size = meta.encoded_size();
577        let encoded_size_u32 = utils::checked_into_u32(meta_encode_size).unwrap_or_else(|_| {
578            panic!(
579                "WARN overflow can't convert meta_encoded_size {} into u32 tables {:?}",
580                meta_encode_size, self.table_ids,
581            )
582        });
583        let meta_offset_u32 = utils::checked_into_u32(meta_offset).unwrap_or_else(|_| {
584            panic!(
585                "WARN overflow can't convert meta_offset {} into u32 tables {:?}",
586                meta_offset, self.table_ids,
587            )
588        });
589        meta.estimated_size = encoded_size_u32
590            .checked_add(meta_offset_u32)
591            .unwrap_or_else(|| {
592                panic!(
593                    "WARN overflow encoded_size_u32 {} meta_offset_u32 {} table_id {:?} table_ids {:?}",
594                    encoded_size_u32, meta_offset_u32, self.last_table_id, self.table_ids
595                )
596            });
597
598        let (avg_key_size, avg_value_size) = if self.table_stats.is_empty() {
599            (0, 0)
600        } else {
601            let total_key_count: usize = self
602                .table_stats
603                .values()
604                .map(|s| s.total_key_count as usize)
605                .sum();
606
607            let total_key_size: usize = self
608                .table_stats
609                .values()
610                .map(|s| s.total_key_size as usize)
611                .sum();
612
613            let total_value_size: usize = self
614                .table_stats
615                .values()
616                .map(|s| s.total_value_size as usize)
617                .sum();
618
619            (
620                total_key_size.checked_div(total_key_count).unwrap_or(0),
621                total_value_size.checked_div(total_key_count).unwrap_or(0),
622            )
623        };
624
625        let (min_epoch, max_epoch) = {
626            if self.epoch_set.is_empty() {
627                (HummockEpoch::MAX, u64::MIN)
628            } else {
629                (
630                    *self.epoch_set.first().unwrap(),
631                    *self.epoch_set.last().unwrap(),
632                )
633            }
634        };
635
636        let vnode_user_key_ranges = self
637            .vnode_range_collector
638            .take()
639            .and_then(|collector| collector.finish(&self.last_full_key));
640
641        let sst_info: SstableInfo = SstableInfoInner {
642            object_id: self.sst_object_id,
643            // use the same sst_id as object_id for initial sst
644            sst_id: self.sst_object_id.as_raw_id().into(),
645            key_range: KeyRange {
646                left: Bytes::from(meta.smallest_key.clone()),
647                right: Bytes::from(meta.largest_key.clone()),
648                right_exclusive,
649            },
650            file_size: meta.estimated_size as u64,
651            table_ids: self.table_ids.into_iter().collect(),
652            meta_offset: meta.meta_offset,
653            stale_key_count,
654            total_key_count,
655            uncompressed_file_size: uncompressed_file_size + meta.encoded_size() as u64,
656            min_epoch,
657            max_epoch,
658            range_tombstone_count: 0,
659            sst_size: meta.estimated_size as u64,
660            filter_type,
661            filter_layout,
662            vnode_statistics: vnode_user_key_ranges,
663        }
664        .into();
665
666        tracing::trace!(
667            "meta_size {} filter_size {} add_key_counts {} stale_key_count {} min_epoch {} max_epoch {} epoch_count {}",
668            meta.encoded_size(),
669            meta.bloom_filter.len(),
670            total_key_count,
671            stale_key_count,
672            min_epoch,
673            max_epoch,
674            self.epoch_set.len()
675        );
676        let filter_size = meta.bloom_filter.len();
677        let sstable_file_size = sst_info.file_size as usize;
678
679        if !meta.block_metas.is_empty() {
680            // fill total_compressed_size
681            let mut last_table_id = meta.block_metas[0].table_id();
682            let mut last_table_stats = self.table_stats.get_mut(&last_table_id).unwrap();
683            for block_meta in &meta.block_metas {
684                let block_table_id = block_meta.table_id();
685                if last_table_id != block_table_id {
686                    last_table_id = block_table_id;
687                    last_table_stats = self.table_stats.get_mut(&last_table_id).unwrap();
688                }
689
690                last_table_stats.total_compressed_size += block_meta.len as u64;
691            }
692        }
693
694        let writer_output = self
695            .writer
696            .finish(meta)
697            .instrument_await("sstable_writer_finish".verbose())
698            .await?;
699        // The timestamp is only used during full GC.
700        //
701        // Ideally object store object's last_modified should be used.
702        // However, it'll incur additional IO overhead since S3 lacks an interface to retrieve the last_modified timestamp after the PUT operation on an object.
703        //
704        // The local timestamp below is expected to precede the last_modified of object store object, given that the object store object is created afterward.
705        // It should help alleviate the clock drift issue.
706
707        let now = SystemTime::now()
708            .duration_since(SystemTime::UNIX_EPOCH)
709            .expect("Clock may have gone backwards")
710            .as_secs();
711        Ok(SstableBuilderOutput::<W::Output> {
712            sst_info: LocalSstableInfo::new(sst_info, self.table_stats, now),
713            writer_output,
714            stats: SstableBuilderOutputStats {
715                filter_size,
716                avg_key_size,
717                avg_value_size,
718                epoch_count: self.epoch_set.len(),
719                block_size_vec: self.block_size_vec,
720                sstable_file_size,
721            },
722        })
723    }
724
725    pub fn approximate_len(&self) -> usize {
726        self.writer.data_len()
727            + self.block_builder.approximate_len()
728            + self.filter_builder.approximate_len()
729    }
730
731    pub async fn build_block(&mut self) -> HummockResult<()> {
732        // Skip empty block.
733        if self.block_builder.is_empty() {
734            return Ok(());
735        }
736
737        let block_meta = self.block_metas.last_mut().unwrap();
738        let uncompressed_block_size = self.block_builder.uncompressed_block_size();
739        block_meta.uncompressed_size = utils::checked_into_u32(uncompressed_block_size)
740            .unwrap_or_else(|_| {
741                panic!(
742                    "WARN overflow can't convert uncompressed_block_size {} into u32 table {:?}",
743                    uncompressed_block_size,
744                    self.block_builder.table_id(),
745                )
746            });
747        let block = self.block_builder.build();
748        self.writer
749            .write_block(block, block_meta)
750            .instrument_await("sstable_write_block".verbose())
751            .await?;
752        self.block_size_vec.push(block.len());
753        let data_len = utils::checked_into_u32(self.writer.data_len()).unwrap_or_else(|_| {
754            panic!(
755                "WARN overflow can't convert writer_data_len {} into u32 table {:?}",
756                self.writer.data_len(),
757                self.block_builder.table_id(),
758            )
759        });
760        block_meta.len = data_len.checked_sub(block_meta.offset).unwrap_or_else(|| {
761            panic!(
762                "data_len should >= meta_offset, found data_len={}, meta_offset={}",
763                data_len, block_meta.offset
764            )
765        });
766
767        self.filter_builder
768            .switch_block(self.memory_limiter.clone());
769
770        if data_len as usize > self.options.capacity * 2 {
771            tracing::warn!(
772                "WARN unexpected block size {} table {:?}",
773                data_len,
774                self.block_builder.table_id()
775            );
776        }
777
778        self.block_builder.clear();
779        Ok(())
780    }
781
782    pub fn is_empty(&self) -> bool {
783        self.writer.data_len() > 0
784    }
785
786    /// Returns true if we roughly reached capacity
787    pub fn reach_capacity(&self) -> bool {
788        self.approximate_len() >= self.options.capacity
789    }
790
791    fn finalize_last_table_stats(&mut self) {
792        if self.table_ids.is_empty() || self.last_table_id.is_none() {
793            return;
794        }
795        self.table_stats.insert(
796            self.last_table_id.unwrap(),
797            std::mem::take(&mut self.last_table_stats),
798        );
799    }
800}
801
802/// Collects vnode key-range hints during SST building.
803struct VnodeUserKeyRangeCollector {
804    max_bytes: usize,
805    current_size: usize,
806    ranges: BTreeMap<VirtualNode, (UserKey<Bytes>, UserKey<Bytes>)>,
807    current_vnode: VirtualNode,
808    range_start_key: Vec<u8>,
809}
810
811impl VnodeUserKeyRangeCollector {
812    fn new(max_bytes: usize) -> Self {
813        Self {
814            max_bytes,
815            current_size: 0,
816            ranges: BTreeMap::new(),
817            current_vnode: VirtualNode::ZERO,
818            range_start_key: Vec::new(),
819        }
820    }
821
822    fn with_limit(max_bytes: Option<usize>) -> Option<Self> {
823        max_bytes.filter(|&n| n > 0).map(Self::new)
824    }
825
826    /// Track vnode boundaries. On vnode switch, seals previous range with `prev_key` as right bound.
827    /// Range: `[first_key_of_vnode, last_key_of_vnode]` (inclusive).
828    fn observe_key(&mut self, vnode: VirtualNode, key: &[u8], prev_key: &[u8]) {
829        if self.current_size >= self.max_bytes {
830            return;
831        }
832
833        // First key
834        if self.range_start_key.is_empty() {
835            self.current_vnode = vnode;
836            self.range_start_key = key.to_vec();
837            return;
838        }
839
840        // Same vnode, nothing to do
841        if vnode == self.current_vnode {
842            return;
843        }
844
845        // Vnode changed: seal previous range
846        self.seal_range(prev_key);
847
848        // Check if budget exhausted after sealing
849        if self.current_size >= self.max_bytes {
850            return;
851        }
852
853        // Start new vnode
854        self.current_vnode = vnode;
855        self.range_start_key = key.to_vec();
856    }
857
858    /// Seal current range. Asserts `vnode/table_id` consistency between left and right keys.
859    fn seal_range(&mut self, right_key: &[u8]) {
860        let left_key = mem::take(&mut self.range_start_key);
861        self.current_size += mem::size_of::<VirtualNode>() + left_key.len() + right_key.len();
862
863        let left_full_key = FullKey::decode(&left_key);
864        let right_full_key = FullKey::decode(right_key);
865        let left_user_key = left_full_key.user_key.copy_into();
866        let right_user_key = right_full_key.user_key.copy_into();
867
868        // Sanity checks for data correctness:
869        // 1. left and right keys have same `vnode`
870        // 2. vnode matches `current_vnode` being tracked
871        // 3. left and right keys have same `table_id`
872        assert_eq!(
873            left_user_key.get_vnode_id(),
874            right_user_key.get_vnode_id(),
875            "vnode changed within range: left_user {:?}, right_user {:?}",
876            left_user_key,
877            right_user_key
878        );
879        assert_eq!(
880            left_user_key.get_vnode_id(),
881            self.current_vnode.to_index(),
882            "vnode mismatch: left {:?}, right {:?}, expected vnode {}",
883            left_user_key,
884            right_user_key,
885            self.current_vnode.to_index()
886        );
887        assert_eq!(
888            left_user_key.table_id, right_user_key.table_id,
889            "table_id changed within range: left {:?}, right {:?}",
890            left_user_key, right_user_key
891        );
892
893        self.ranges
894            .insert(self.current_vnode, (left_user_key, right_user_key));
895    }
896
897    /// Returns `Some` if >1 vnodes collected, `None` otherwise (single-vnode needs no hints).
898    fn finish(mut self, last_key: &[u8]) -> Option<VnodeStatistics> {
899        if !self.range_start_key.is_empty() {
900            self.seal_range(last_key);
901        }
902
903        if self.ranges.len() > 1 {
904            // Validate all ranges belong to the same table_id before building VnodeStatistics
905            let mut table_ids = self
906                .ranges
907                .values()
908                .flat_map(|(left, right)| [left.table_id, right.table_id]);
909            if let Some(first_table_id) = table_ids.next() {
910                for table_id in table_ids {
911                    assert_eq!(
912                        table_id, first_table_id,
913                        "all vnode ranges must belong to the same table_id, found {:?} and {:?}",
914                        table_id, first_table_id
915                    );
916                }
917            }
918
919            Some(VnodeStatistics::from_map(self.ranges))
920        } else {
921            None
922        }
923    }
924}
925
926pub struct SstableBuilderOutputStats {
927    filter_size: usize,
928    avg_key_size: usize,
929    avg_value_size: usize,
930    epoch_count: usize,
931    block_size_vec: Vec<usize>, // for statistics
932    sstable_file_size: usize,
933}
934
935impl SstableBuilderOutputStats {
936    pub fn report_stats(&self, metrics: &Arc<CompactorMetrics>) {
937        if self.filter_size != 0 {
938            metrics
939                .sstable_bloom_filter_size
940                .observe(self.filter_size as _);
941        }
942
943        if self.sstable_file_size != 0 {
944            metrics
945                .sstable_file_size
946                .observe(self.sstable_file_size as _);
947        }
948
949        if self.avg_key_size != 0 {
950            metrics.sstable_avg_key_size.observe(self.avg_key_size as _);
951        }
952
953        if self.avg_value_size != 0 {
954            metrics
955                .sstable_avg_value_size
956                .observe(self.avg_value_size as _);
957        }
958
959        if self.epoch_count != 0 {
960            metrics
961                .sstable_distinct_epoch_count
962                .observe(self.epoch_count as _);
963        }
964
965        if !self.block_size_vec.is_empty() {
966            for block_size in &self.block_size_vec {
967                metrics.sstable_block_size.observe(*block_size as _);
968            }
969        }
970    }
971}
972
973#[cfg(test)]
974pub(super) mod tests {
975    use std::collections::{Bound, HashMap};
976
977    use risingwave_common::catalog::TableId;
978    use risingwave_common::hash::VirtualNode;
979    use risingwave_common::util::epoch::test_epoch;
980    use risingwave_hummock_sdk::key::UserKey;
981
982    use super::*;
983    use crate::assert_bytes_eq;
984    use crate::compaction_catalog_manager::{
985        CompactionCatalogAgent, DummyFilterKeyExtractor, MultiFilterKeyExtractor,
986    };
987    use crate::hummock::iterator::test_utils::mock_sstable_store;
988    use crate::hummock::sstable::xor_filter::BlockedXor16FilterBuilder;
989    use crate::hummock::test_utils::{
990        TEST_KEYS_COUNT, default_builder_opt_for_test, gen_test_sstable_impl, mock_sst_writer,
991        test_key_of, test_value_of,
992    };
993    use crate::hummock::{CachePolicy, Sstable, SstableWriterOptions, Xor8FilterBuilder};
994    use crate::monitor::StoreLocalStatistic;
995
996    #[tokio::test]
997    async fn test_empty() {
998        let opt = SstableBuilderOptions {
999            capacity: 0,
1000            block_capacity: 4096,
1001            restart_interval: 16,
1002            bloom_false_positive: 0.001,
1003            ..Default::default()
1004        };
1005
1006        let table_id_to_vnode = HashMap::from_iter(vec![(0, VirtualNode::COUNT_FOR_TEST)]);
1007        let table_id_to_watermark_serde = HashMap::from_iter(vec![(0, None)]);
1008        let b = SstableBuilder::for_test(
1009            0,
1010            mock_sst_writer(&opt),
1011            opt,
1012            table_id_to_vnode,
1013            table_id_to_watermark_serde,
1014        );
1015
1016        b.finish().await.unwrap();
1017    }
1018
1019    fn encode_full_key(vnode: VirtualNode, table_key_suffix: &[u8]) -> Vec<u8> {
1020        let mut table_key = vnode.to_be_bytes().to_vec();
1021        table_key.extend_from_slice(table_key_suffix);
1022        FullKey::for_test(TableId::default(), table_key, 0).encode()
1023    }
1024
1025    fn table_key_of(vnode: VirtualNode, suffix: &[u8]) -> Vec<u8> {
1026        let mut key = vnode.to_be_bytes().to_vec();
1027        key.extend_from_slice(suffix);
1028        key
1029    }
1030
1031    #[test]
1032    fn test_vnode_user_key_range_basic_collection() {
1033        // Test basic multi-vnode collection with boundary semantics verification.
1034        // Validates: vnode switching triggers range sealing, boundaries are inclusive (right_exclusive=false).
1035        let mut collector = VnodeUserKeyRangeCollector::with_limit(Some(1024)).unwrap();
1036        let vnode_1 = VirtualNode::from_index(1);
1037        let vnode_2 = VirtualNode::from_index(2);
1038
1039        let k1 = encode_full_key(vnode_1, b"k1");
1040        let k2 = encode_full_key(vnode_1, b"k2");
1041        let k3 = encode_full_key(vnode_2, b"k3");
1042        let k4 = encode_full_key(vnode_2, b"k4");
1043
1044        collector.observe_key(vnode_1, &k1, &[]);
1045        collector.observe_key(vnode_1, &k2, &k1);
1046        collector.observe_key(vnode_2, &k3, &k2);
1047        collector.observe_key(vnode_2, &k4, &k3);
1048
1049        let info = collector.finish(&k4).unwrap();
1050        assert_eq!(info.vnode_user_key_ranges().len(), 2);
1051
1052        // Verify vnode_1: left = first key, right = last key before switch
1053        let (range1_left, range1_right) = info.get_vnode_user_key_range(vnode_1).unwrap();
1054        assert_eq!(range1_left.table_key.as_ref(), table_key_of(vnode_1, b"k1"));
1055        assert_eq!(
1056            range1_right.table_key.as_ref(),
1057            table_key_of(vnode_1, b"k2")
1058        );
1059
1060        // Verify vnode_2: left = first key, right = SST's last key
1061        let (range2_left, range2_right) = info.get_vnode_user_key_range(vnode_2).unwrap();
1062        assert_eq!(range2_left.table_key.as_ref(), table_key_of(vnode_2, b"k3"));
1063        assert_eq!(
1064            range2_right.table_key.as_ref(),
1065            table_key_of(vnode_2, b"k4")
1066        );
1067    }
1068
1069    #[test]
1070    fn test_vnode_user_key_range_capacity_limit() {
1071        // Test "allow over-limit write" semantics: inserting a range may exceed capacity,
1072        // but stops before starting the next range if it would exceed the limit.
1073        //
1074        // Calculation based on encoded key sizes:
1075        // Each range = sizeof(VirtualNode) + left.len() + right.len().
1076        // Use a limit that:
1077        //   - allows writing vnode_1 (range_size),
1078        //   - allows over-limit write of vnode_2 (range_size + estimated_next),
1079        //   - stops before vnode_3 (2 * range_size + estimated_next > limit).
1080        let vnode_1 = VirtualNode::from_index(1);
1081        let vnode_2 = VirtualNode::from_index(2);
1082        let vnode_3 = VirtualNode::from_index(3);
1083        let vnode_4 = VirtualNode::from_index(4);
1084
1085        let k1 = encode_full_key(vnode_1, b"k1");
1086        let k2 = encode_full_key(vnode_1, b"k2");
1087        let k3 = encode_full_key(vnode_2, b"k3");
1088        let k4 = encode_full_key(vnode_2, b"k4");
1089        let k5 = encode_full_key(vnode_3, b"k5");
1090        let k6 = encode_full_key(vnode_3, b"k6");
1091        let k7 = encode_full_key(vnode_4, b"k7");
1092
1093        let range_size = mem::size_of::<VirtualNode>() + k1.len() + k2.len();
1094        let estimated_next_size = mem::size_of::<VirtualNode>() + k3.len();
1095        let limit = range_size + estimated_next_size; // allow second range, stop before third
1096        let mut collector = VnodeUserKeyRangeCollector::with_limit(Some(limit)).unwrap();
1097
1098        collector.observe_key(vnode_1, &k1, &[]);
1099        collector.observe_key(vnode_1, &k2, &k1);
1100        collector.observe_key(vnode_2, &k3, &k2); // Seals vnode_1 (6 bytes)
1101        collector.observe_key(vnode_2, &k4, &k3);
1102        collector.observe_key(vnode_3, &k5, &k4); // Seals vnode_2 (12 bytes total), stops before vnode_3
1103        collector.observe_key(vnode_3, &k6, &k5); // Ignored
1104        collector.observe_key(vnode_4, &k7, &k6); // Ignored
1105
1106        let info = collector.finish(&k7).unwrap();
1107        assert_eq!(
1108            info.vnode_user_key_ranges().len(),
1109            2,
1110            "Should collect exactly 2 vnodes"
1111        );
1112
1113        // Verify collected vnodes have correct boundaries
1114        let (range1_left, range1_right) = info.get_vnode_user_key_range(vnode_1).unwrap();
1115        assert_eq!(range1_left.table_key.as_ref(), table_key_of(vnode_1, b"k1"));
1116        assert_eq!(
1117            range1_right.table_key.as_ref(),
1118            table_key_of(vnode_1, b"k2")
1119        );
1120
1121        let (range2_left, range2_right) = info.get_vnode_user_key_range(vnode_2).unwrap();
1122        assert_eq!(range2_left.table_key.as_ref(), table_key_of(vnode_2, b"k3"));
1123        assert_eq!(
1124            range2_right.table_key.as_ref(),
1125            table_key_of(vnode_2, b"k4")
1126        );
1127
1128        // Verify stopped vnodes are not collected
1129        assert!(info.get_vnode_user_key_range(vnode_3).is_none());
1130        assert!(info.get_vnode_user_key_range(vnode_4).is_none());
1131    }
1132
1133    #[test]
1134    fn test_vnode_user_key_range_sparse_distribution() {
1135        // Test non-consecutive vnodes (production scenario: hash-based data distribution).
1136        // Validates: collector handles sparse vnode indices correctly, gaps are not filled.
1137        let mut collector = VnodeUserKeyRangeCollector::with_limit(Some(1024)).unwrap();
1138        let vnode_5 = VirtualNode::from_index(5);
1139        let vnode_10 = VirtualNode::from_index(10);
1140        let vnode_100 = VirtualNode::from_index(100);
1141
1142        let keys = vec![
1143            (vnode_5, encode_full_key(vnode_5, b"key_005_001")),
1144            (vnode_5, encode_full_key(vnode_5, b"key_005_002")),
1145            (vnode_5, encode_full_key(vnode_5, b"key_005_999")),
1146            (vnode_10, encode_full_key(vnode_10, b"key_010_001")),
1147            (vnode_10, encode_full_key(vnode_10, b"key_010_100")),
1148            (vnode_100, encode_full_key(vnode_100, b"key_100_001")),
1149            (vnode_100, encode_full_key(vnode_100, b"key_100_999")),
1150        ];
1151
1152        let mut prev_key = Vec::new();
1153        for (vnode, key) in &keys {
1154            collector.observe_key(*vnode, key, &prev_key);
1155            prev_key = key.clone();
1156        }
1157
1158        let info = collector.finish(&prev_key).unwrap();
1159        assert_eq!(info.vnode_user_key_ranges().len(), 3);
1160
1161        // Verify each collected vnode has correct boundaries
1162        let (range5_left, range5_right) = info.get_vnode_user_key_range(vnode_5).unwrap();
1163        assert_eq!(
1164            range5_left.table_key.as_ref(),
1165            table_key_of(vnode_5, b"key_005_001")
1166        );
1167        assert_eq!(
1168            range5_right.table_key.as_ref(),
1169            table_key_of(vnode_5, b"key_005_999")
1170        );
1171
1172        let (range10_left, range10_right) = info.get_vnode_user_key_range(vnode_10).unwrap();
1173        assert_eq!(
1174            range10_left.table_key.as_ref(),
1175            table_key_of(vnode_10, b"key_010_001")
1176        );
1177        assert_eq!(
1178            range10_right.table_key.as_ref(),
1179            table_key_of(vnode_10, b"key_010_100")
1180        );
1181
1182        let (range100_left, range100_right) = info.get_vnode_user_key_range(vnode_100).unwrap();
1183        assert_eq!(
1184            range100_left.table_key.as_ref(),
1185            table_key_of(vnode_100, b"key_100_001")
1186        );
1187        assert_eq!(
1188            range100_right.table_key.as_ref(),
1189            table_key_of(vnode_100, b"key_100_999")
1190        );
1191
1192        // Verify gaps are not filled (no data for vnodes 0, 7, 50)
1193        assert!(
1194            info.get_vnode_user_key_range(VirtualNode::from_index(0))
1195                .is_none()
1196        );
1197        assert!(
1198            info.get_vnode_user_key_range(VirtualNode::from_index(7))
1199                .is_none()
1200        );
1201        assert!(
1202            info.get_vnode_user_key_range(VirtualNode::from_index(50))
1203                .is_none()
1204        );
1205    }
1206
1207    #[test]
1208    fn test_vnode_user_key_range_edge_cases() {
1209        // Test 1: Configuration disabled (None or 0) should return None
1210        assert!(VnodeUserKeyRangeCollector::with_limit(None).is_none());
1211        assert!(VnodeUserKeyRangeCollector::with_limit(Some(0)).is_none());
1212
1213        // Test 2: Empty collector (no keys) should return None
1214        let collector = VnodeUserKeyRangeCollector::with_limit(Some(1024)).unwrap();
1215        assert!(
1216            collector
1217                .finish(&encode_full_key(VirtualNode::ZERO, b"any"))
1218                .is_none()
1219        );
1220
1221        // Test 3: Single vnode (optimization: don't emit hints for single-vnode SSTs)
1222        let vnode = VirtualNode::from_index(5);
1223        let mut collector = VnodeUserKeyRangeCollector::with_limit(Some(1024)).unwrap();
1224        let a = encode_full_key(vnode, b"a");
1225        let b = encode_full_key(vnode, b"b");
1226        let c = encode_full_key(vnode, b"c");
1227        collector.observe_key(vnode, &a, &[]);
1228        collector.observe_key(vnode, &b, &a);
1229        collector.observe_key(vnode, &c, &b);
1230        assert!(
1231            collector.finish(&c).is_none(),
1232            "Single-vnode SST should not emit hints"
1233        );
1234
1235        // Test 4: Extreme limit (1 byte) - first vnode exceeds, becomes single-vnode
1236        // Range size = sizeof(VirtualNode=2) + encoded_key_len * 2
1237        // With 1-byte limit: vnode_1 range exceeds but allowed (over-limit write)
1238        // vnode_2 check: current_size >= max_bytes, stop. Result: only vnode_1, returns None (single-vnode)
1239        let mut collector = VnodeUserKeyRangeCollector::with_limit(Some(1)).unwrap();
1240        let vnode_1 = VirtualNode::from_index(1);
1241        let vnode_2 = VirtualNode::from_index(2);
1242        let key1 = encode_full_key(vnode_1, b"key1");
1243        let key2 = encode_full_key(vnode_1, b"key2");
1244        let key3 = encode_full_key(vnode_2, b"key3");
1245        collector.observe_key(vnode_1, &key1, &[]);
1246        collector.observe_key(vnode_1, &key2, &key1);
1247        collector.observe_key(vnode_2, &key3, &key2); // Seals vnode_1, stops before vnode_2
1248        assert!(
1249            collector.finish(&key3).is_none(),
1250            "Only 1 vnode collected, should return None"
1251        );
1252    }
1253
1254    #[tokio::test]
1255    async fn test_add_raw_block_compression_and_fallback() {
1256        let new_builder = |compression_algorithm| {
1257            let options = SstableBuilderOptions {
1258                block_capacity: 2048,
1259                compression_algorithm,
1260                ..Default::default()
1261            };
1262            SstableBuilder::<_, BlockedXor16FilterBuilder>::new(
1263                1,
1264                mock_sst_writer(&options),
1265                BlockedXor16FilterBuilder::create(options.filter_builder_options()),
1266                options,
1267                CompactionCatalogAgent::for_test(vec![0]),
1268                None,
1269            )
1270        };
1271        let value = vec![7; 600];
1272        let mut source = new_builder(CompressionAlgorithm::None);
1273        for index in [2, 3] {
1274            source
1275                .add(test_key_of(index).to_ref(), HummockValue::put(&value))
1276                .await
1277                .unwrap();
1278        }
1279        let (data, meta) = source.finish().await.unwrap().writer_output;
1280        let source = Sstable::new(1.into(), meta, false);
1281        assert_eq!(source.meta.block_metas.len(), 1);
1282        let source_meta = &source.meta.block_metas[0];
1283        let raw = data.slice(..source_meta.len as usize);
1284
1285        for compression in [CompressionAlgorithm::Lz4, CompressionAlgorithm::Zstd] {
1286            for prefix_size in [16, 600] {
1287                let mut builder = new_builder(compression);
1288                let prefix_value = vec![1; prefix_size];
1289                builder
1290                    .add(test_key_of(1).to_ref(), HummockValue::put(&prefix_value))
1291                    .await
1292                    .unwrap();
1293                let written = builder
1294                    .add_raw_block(
1295                        raw.clone(),
1296                        source.filter_reader.get_block_raw_filter(0),
1297                        test_key_of(2),
1298                        test_key_of(3).encode(),
1299                        source_meta.clone(),
1300                    )
1301                    .await
1302                    .unwrap();
1303                let (data, meta) = builder.finish().await.unwrap().writer_output;
1304                if prefix_size == 16 {
1305                    // The pending block is below 2048 / 4, so it absorbs the raw input.
1306                    assert_eq!(written, None);
1307                    assert_eq!(meta.block_metas.len(), 1);
1308                } else {
1309                    assert_eq!(meta.block_metas.len(), 2);
1310                    let copied = &meta.block_metas[1];
1311                    assert_eq!(written, Some(copied.len as usize));
1312                    assert!(copied.len < source_meta.len);
1313                    assert_eq!(copied.uncompressed_size, source_meta.uncompressed_size);
1314                }
1315                let mut expected_index = 1;
1316                let mut offset = 0;
1317                for block_meta in &meta.block_metas {
1318                    assert_eq!(block_meta.offset, offset);
1319                    offset += block_meta.len;
1320                    let block = data.slice(block_meta.offset as usize..offset as usize);
1321                    assert_eq!(Block::get_algorithm(&block).unwrap(), compression);
1322                    let block =
1323                        Block::decode(block, block_meta.uncompressed_size as usize).unwrap();
1324                    let mut iter =
1325                        BlockIterator::new(BlockHolder::from_owned_block(Box::new(block)));
1326                    iter.seek_to_first();
1327                    while iter.is_valid() {
1328                        assert_eq!(iter.key(), test_key_of(expected_index).to_ref());
1329                        let expected_value = if expected_index == 1 {
1330                            &prefix_value
1331                        } else {
1332                            &value
1333                        };
1334                        assert_eq!(
1335                            HummockValue::from_slice(iter.value()).unwrap(),
1336                            HummockValue::put(expected_value.as_slice())
1337                        );
1338                        expected_index += 1;
1339                        iter.next();
1340                    }
1341                }
1342                assert_eq!(expected_index, 4);
1343            }
1344        }
1345    }
1346
1347    #[tokio::test]
1348    async fn test_basic() {
1349        let opt = default_builder_opt_for_test();
1350
1351        let table_id_to_vnode = HashMap::from_iter(vec![(0, VirtualNode::COUNT_FOR_TEST)]);
1352        let table_id_to_watermark_serde = HashMap::from_iter(vec![(0, None)]);
1353        let mut b = SstableBuilder::for_test(
1354            0,
1355            mock_sst_writer(&opt),
1356            opt,
1357            table_id_to_vnode,
1358            table_id_to_watermark_serde,
1359        );
1360
1361        for i in 0..TEST_KEYS_COUNT {
1362            b.add_for_test(
1363                test_key_of(i).to_ref(),
1364                HummockValue::put(&test_value_of(i)),
1365            )
1366            .await
1367            .unwrap();
1368        }
1369
1370        let output = b.finish().await.unwrap();
1371        let info = output.sst_info.sst_info;
1372
1373        assert_bytes_eq!(test_key_of(0).encode(), info.key_range.left);
1374        assert_bytes_eq!(
1375            test_key_of(TEST_KEYS_COUNT - 1).encode(),
1376            info.key_range.right
1377        );
1378        let (data, meta) = output.writer_output;
1379        assert_eq!(info.file_size, meta.estimated_size as u64);
1380        let offset = info.meta_offset as usize;
1381        let meta2 = SstableMeta::decode(&data[offset..]).unwrap();
1382        assert_eq!(meta2, meta);
1383    }
1384
1385    async fn test_with_xor_filter_builder<F: FilterBuilder>(
1386        bloom_false_positive: f64,
1387        expected_filter_type: PbSstableFilterType,
1388        expected_filter_layout: PbSstableFilterLayout,
1389    ) {
1390        let key_count = 1000;
1391
1392        let opts = SstableBuilderOptions {
1393            capacity: 0,
1394            block_capacity: 4096,
1395            restart_interval: 16,
1396            bloom_false_positive,
1397            ..Default::default()
1398        };
1399
1400        // build remote table
1401        let sstable_store = mock_sstable_store().await;
1402        let table_id_to_vnode = HashMap::from_iter(vec![(0, VirtualNode::COUNT_FOR_TEST)]);
1403        let table_id_to_watermark_serde = HashMap::from_iter(vec![(0, None)]);
1404        let sst_info = gen_test_sstable_impl::<Vec<u8>, F>(
1405            opts,
1406            0,
1407            (0..TEST_KEYS_COUNT).map(|i| (test_key_of(i), HummockValue::put(test_value_of(i)))),
1408            sstable_store.clone(),
1409            CachePolicy::NotFill,
1410            table_id_to_vnode,
1411            table_id_to_watermark_serde,
1412        )
1413        .await;
1414        assert_eq!(sst_info.filter_type, expected_filter_type);
1415        assert_eq!(sst_info.filter_layout, expected_filter_layout);
1416        let table = sstable_store
1417            .sstable(&sst_info, &mut StoreLocalStatistic::default())
1418            .await
1419            .unwrap();
1420
1421        assert!(table.has_filter());
1422        for i in 0..key_count {
1423            let full_key = test_key_of(i);
1424            let hash = Sstable::hash_for_filter(full_key.user_key.encode().as_slice(), 0);
1425            let key_ref = full_key.user_key.as_ref();
1426            assert!(
1427                table.may_match_hash(&(Bound::Included(key_ref), Bound::Included(key_ref)), hash),
1428                "failed at {}",
1429                i
1430            );
1431        }
1432    }
1433
1434    #[tokio::test]
1435    async fn test_xor_filter_builder_output() {
1436        test_with_xor_filter_builder::<Xor16FilterBuilder>(
1437            0.0,
1438            PbSstableFilterType::SstableFilterXor16,
1439            PbSstableFilterLayout::Plain,
1440        )
1441        .await;
1442        test_with_xor_filter_builder::<Xor16FilterBuilder>(
1443            0.01,
1444            PbSstableFilterType::SstableFilterXor16,
1445            PbSstableFilterLayout::Plain,
1446        )
1447        .await;
1448        test_with_xor_filter_builder::<Xor8FilterBuilder>(
1449            0.01,
1450            PbSstableFilterType::SstableFilterXor8,
1451            PbSstableFilterLayout::Plain,
1452        )
1453        .await;
1454        test_with_xor_filter_builder::<BlockedXor16FilterBuilder>(
1455            0.01,
1456            PbSstableFilterType::SstableFilterXor16,
1457            PbSstableFilterLayout::Blocked,
1458        )
1459        .await;
1460    }
1461
1462    #[tokio::test]
1463    async fn test_no_xor_filter_block() {
1464        let opts = SstableBuilderOptions::default();
1465        // build remote table
1466        let sstable_store = mock_sstable_store().await;
1467        let writer_opts = SstableWriterOptions::default();
1468        let object_id = 1;
1469        let writer = sstable_store
1470            .clone()
1471            .create_sst_writer(object_id, writer_opts);
1472        let mut filter = MultiFilterKeyExtractor::default();
1473        filter.register(
1474            1.into(),
1475            FilterKeyExtractorImpl::Dummy(DummyFilterKeyExtractor),
1476        );
1477        filter.register(
1478            2.into(),
1479            FilterKeyExtractorImpl::FullKey(FullKeyFilterKeyExtractor),
1480        );
1481        filter.register(
1482            3.into(),
1483            FilterKeyExtractorImpl::Dummy(DummyFilterKeyExtractor),
1484        );
1485
1486        let table_id_to_vnode = HashMap::from_iter(vec![
1487            (1.into(), VirtualNode::COUNT_FOR_TEST),
1488            (2.into(), VirtualNode::COUNT_FOR_TEST),
1489            (3.into(), VirtualNode::COUNT_FOR_TEST),
1490        ]);
1491        let table_id_to_watermark_serde =
1492            HashMap::from_iter(vec![(1.into(), None), (2.into(), None), (3.into(), None)]);
1493
1494        let compaction_catalog_agent_ref = Arc::new(CompactionCatalogAgent::new(
1495            FilterKeyExtractorImpl::Multi(filter),
1496            table_id_to_vnode,
1497            table_id_to_watermark_serde,
1498            HashMap::default(),
1499        ));
1500
1501        let mut builder = SstableBuilder::new(
1502            object_id,
1503            writer,
1504            BlockedXor16FilterBuilder::create(opts.filter_builder_options()),
1505            opts,
1506            compaction_catalog_agent_ref,
1507            None,
1508        );
1509
1510        let key_count: usize = 10000;
1511        for table_id in 1..4 {
1512            let mut table_key = VirtualNode::ZERO.to_be_bytes().to_vec();
1513            for idx in 0..key_count {
1514                table_key.resize(VirtualNode::SIZE, 0);
1515                table_key.extend_from_slice(format!("key_test_{:05}", idx * 2).as_bytes());
1516                let k = UserKey::for_test(TableId::new(table_id), table_key.as_ref());
1517                let v = test_value_of(idx);
1518                builder
1519                    .add(
1520                        FullKey::from_user_key(k, test_epoch(1)),
1521                        HummockValue::put(v.as_ref()),
1522                    )
1523                    .await
1524                    .unwrap();
1525            }
1526        }
1527        let ret = builder.finish().await.unwrap();
1528        let sst_info = ret.sst_info.sst_info.clone();
1529        ret.writer_output.await.unwrap().unwrap();
1530        let table = sstable_store
1531            .sstable(&sst_info, &mut StoreLocalStatistic::default())
1532            .await
1533            .unwrap();
1534        let mut table_key = VirtualNode::ZERO.to_be_bytes().to_vec();
1535        for idx in 0..key_count {
1536            table_key.resize(VirtualNode::SIZE, 0);
1537            table_key.extend_from_slice(format!("key_test_{:05}", idx * 2).as_bytes());
1538            let k = UserKey::for_test(TableId::new(2), table_key.as_slice());
1539            let hash = Sstable::hash_for_filter(&k.encode(), 2);
1540            let key_ref = k.as_ref();
1541            assert!(
1542                table.may_match_hash(&(Bound::Included(key_ref), Bound::Included(key_ref)), hash)
1543            );
1544        }
1545    }
1546}