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