Skip to main content

risingwave_hummock_sdk/compaction_group/
hummock_version_ext.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::borrow::Borrow;
16use std::cmp::Ordering;
17use std::collections::hash_map::Entry;
18use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet};
19use std::iter::once;
20use std::sync::{Arc, LazyLock};
21
22use bytes::Bytes;
23use itertools::Itertools;
24use risingwave_common::catalog::TableId;
25use risingwave_common::hash::VnodeBitmapExt;
26use risingwave_common::log::LogSuppressor;
27use risingwave_pb::hummock::{
28    CompactionConfig, CompatibilityVersion, PbLevelType, StateTableInfo, StateTableInfoDelta,
29};
30use tracing::warn;
31
32use super::group_split::split_sst_with_table_ids;
33use super::{StateTableId, group_split};
34use crate::change_log::{EpochNewChangeLogCommon, TableChangeLogCommon};
35use crate::compact_task::is_compaction_task_expired;
36use crate::compaction_group::StaticCompactionGroupId;
37use crate::key_range::KeyRangeCommon;
38use crate::level::{Level, LevelCommon, Levels, OverlappingLevel};
39use crate::sstable_info::SstableInfo;
40use crate::table_watermark::{ReadTableWatermark, TableWatermarks};
41use crate::vector_index::apply_vector_index_delta;
42use crate::version::{
43    GroupDelta, GroupDeltaCommon, HummockVersion, HummockVersionCommon, HummockVersionDeltaCommon,
44    IntraLevelDelta, IntraLevelDeltaCommon, ObjectIdReader, SstableIdReader,
45};
46use crate::{
47    CompactionGroupId, HummockObjectId, HummockSstableId, HummockSstableObjectId, can_concat,
48};
49
50#[derive(Debug, Clone, Default)]
51pub struct SstDeltaInfo {
52    pub insert_sst_level: u32,
53    pub insert_sst_infos: Vec<SstableInfo>,
54    pub delete_sst_object_ids: Vec<HummockSstableObjectId>,
55}
56
57pub type BranchedSstInfo = HashMap<CompactionGroupId, Vec<HummockSstableId>>;
58
59impl HummockVersionCommon<SstableInfo> {
60    pub fn get_compaction_group_levels(&self, compaction_group_id: CompactionGroupId) -> &Levels {
61        self.levels
62            .get(&compaction_group_id)
63            .unwrap_or_else(|| panic!("compaction group {} does not exist", compaction_group_id))
64    }
65
66    pub fn get_compaction_group_levels_mut(
67        &mut self,
68        compaction_group_id: CompactionGroupId,
69    ) -> &mut Levels {
70        self.levels
71            .get_mut(&compaction_group_id)
72            .unwrap_or_else(|| panic!("compaction group {} does not exist", compaction_group_id))
73    }
74
75    // Only scan SST infos from levels in the specified compaction group.
76    pub fn get_sst_ids_by_group_id(
77        &self,
78        compaction_group_id: CompactionGroupId,
79    ) -> impl Iterator<Item = HummockSstableId> + '_ {
80        self.levels
81            .get(&compaction_group_id)
82            .into_iter()
83            .flat_map(|level| level.l0.sub_levels.iter().rev().chain(level.levels.iter()))
84            .flat_map(|level| level.table_infos.iter())
85            .map(|s| s.sst_id)
86    }
87
88    /// Prune stale table ids that no longer exist in `state_table_info` from SST metadata.
89    ///
90    /// This is used to normalize recovered versions from old clusters where dropped table ids
91    /// may still exist in persisted SST metadata.
92    pub fn prune_stale_table_ids_from_ssts(&mut self) -> usize {
93        let live_table_ids: HashSet<_> = self.state_table_info.info().keys().copied().collect();
94        // Older checkpoints may rely on deprecated `member_table_ids` before `state_table_info`
95        // is backfilled.
96        if live_table_ids.is_empty()
97            && self.levels.values().any(|levels| {
98                #[expect(deprecated)]
99                {
100                    !levels.member_table_ids.is_empty()
101                }
102            })
103        {
104            return 0;
105        }
106
107        let mut pruned_table_ids = HashSet::new();
108
109        for levels in self.levels.values_mut() {
110            let stale_table_ids = levels
111                .l0
112                .sub_levels
113                .iter()
114                .chain(levels.levels.iter())
115                .flat_map(|level| level.table_infos.iter())
116                .flat_map(|sst| sst.table_ids.iter().copied())
117                .filter(|table_id| !live_table_ids.contains(table_id))
118                .collect::<HashSet<_>>();
119
120            if stale_table_ids.is_empty() {
121                continue;
122            }
123
124            pruned_table_ids.extend(stale_table_ids.iter().copied());
125            levels.prune_table_ids_from_ssts(&stale_table_ids);
126        }
127
128        pruned_table_ids.len()
129    }
130
131    pub fn level_iter<F: FnMut(&Level) -> bool>(
132        &self,
133        compaction_group_id: CompactionGroupId,
134        mut f: F,
135    ) {
136        if let Some(levels) = self.levels.get(&compaction_group_id) {
137            for sub_level in &levels.l0.sub_levels {
138                if !f(sub_level) {
139                    return;
140                }
141            }
142            for level in &levels.levels {
143                if !f(level) {
144                    return;
145                }
146            }
147        }
148    }
149
150    pub fn num_levels(&self, compaction_group_id: CompactionGroupId) -> usize {
151        // l0 is currently separated from all levels
152        self.levels
153            .get(&compaction_group_id)
154            .map(|group| group.levels.len() + 1)
155            .unwrap_or(0)
156    }
157
158    pub fn safe_epoch_table_watermarks(
159        &self,
160        existing_table_ids: &[TableId],
161    ) -> BTreeMap<TableId, TableWatermarks> {
162        safe_epoch_table_watermarks_impl(&self.table_watermarks, existing_table_ids)
163    }
164}
165
166pub fn safe_epoch_table_watermarks_impl(
167    table_watermarks: &HashMap<TableId, Arc<TableWatermarks>>,
168    existing_table_ids: &[TableId],
169) -> BTreeMap<TableId, TableWatermarks> {
170    fn extract_single_table_watermark(
171        table_watermarks: &TableWatermarks,
172    ) -> Option<TableWatermarks> {
173        if let Some((first_epoch, first_epoch_watermark)) = table_watermarks.watermarks.first() {
174            Some(TableWatermarks {
175                watermarks: vec![(*first_epoch, first_epoch_watermark.clone())],
176                direction: table_watermarks.direction,
177                watermark_type: table_watermarks.watermark_type,
178            })
179        } else {
180            None
181        }
182    }
183    table_watermarks
184        .iter()
185        .filter_map(|(table_id, table_watermarks)| {
186            if !existing_table_ids.contains(table_id) {
187                None
188            } else {
189                extract_single_table_watermark(table_watermarks)
190                    .map(|table_watermarks| (*table_id, table_watermarks))
191            }
192        })
193        .collect()
194}
195
196pub fn safe_epoch_read_table_watermarks_impl(
197    safe_epoch_watermarks: BTreeMap<TableId, TableWatermarks>,
198) -> BTreeMap<TableId, ReadTableWatermark> {
199    safe_epoch_watermarks
200        .into_iter()
201        .map(|(table_id, watermarks)| {
202            assert_eq!(watermarks.watermarks.len(), 1);
203            let vnode_watermarks = &watermarks.watermarks.first().expect("should exist").1;
204            let mut vnode_watermark_map = BTreeMap::new();
205            for vnode_watermark in vnode_watermarks.iter() {
206                let watermark = Bytes::copy_from_slice(vnode_watermark.watermark());
207                for vnode in vnode_watermark.vnode_bitmap().iter_vnodes() {
208                    assert!(
209                        vnode_watermark_map
210                            .insert(vnode, watermark.clone())
211                            .is_none(),
212                        "duplicate table watermark on vnode {}",
213                        vnode.to_index()
214                    );
215                }
216            }
217            (
218                table_id,
219                ReadTableWatermark {
220                    direction: watermarks.direction,
221                    vnode_watermarks: vnode_watermark_map,
222                },
223            )
224        })
225        .collect()
226}
227
228impl HummockVersionCommon<SstableInfo> {
229    pub fn count_new_ssts_in_group_split(
230        &self,
231        parent_group_id: CompactionGroupId,
232        split_key: Bytes,
233    ) -> u64 {
234        self.levels
235            .get(&parent_group_id)
236            .map_or(0, |parent_levels| {
237                let l0 = &parent_levels.l0;
238                let mut split_count = 0;
239                for sub_level in &l0.sub_levels {
240                    assert!(!sub_level.table_infos.is_empty());
241
242                    if sub_level.level_type == PbLevelType::Overlapping {
243                        // TODO: use table_id / vnode / key_range filter
244                        split_count += sub_level
245                            .table_infos
246                            .iter()
247                            .map(|sst| {
248                                if let group_split::SstSplitType::Both =
249                                    group_split::need_to_split(sst, split_key.clone())
250                                {
251                                    2
252                                } else {
253                                    0
254                                }
255                            })
256                            .sum::<u64>();
257                        continue;
258                    }
259
260                    let pos = group_split::get_split_pos(&sub_level.table_infos, split_key.clone());
261                    let sst = sub_level.table_infos.get(pos).unwrap();
262
263                    if let group_split::SstSplitType::Both =
264                        group_split::need_to_split(sst, split_key.clone())
265                    {
266                        split_count += 2;
267                    }
268                }
269
270                for level in &parent_levels.levels {
271                    if level.table_infos.is_empty() {
272                        continue;
273                    }
274                    let pos = group_split::get_split_pos(&level.table_infos, split_key.clone());
275                    let sst = level.table_infos.get(pos).unwrap();
276                    if let group_split::SstSplitType::Both =
277                        group_split::need_to_split(sst, split_key.clone())
278                    {
279                        split_count += 2;
280                    }
281                }
282
283                split_count
284            })
285    }
286
287    pub fn init_with_parent_group(
288        &mut self,
289        parent_group_id: CompactionGroupId,
290        group_id: CompactionGroupId,
291        member_table_ids: BTreeSet<StateTableId>,
292        new_sst_start_id: HummockSstableId,
293    ) {
294        let mut new_sst_id = new_sst_start_id;
295        if parent_group_id == StaticCompactionGroupId::NewCompactionGroup {
296            if new_sst_start_id != 0 {
297                if cfg!(debug_assertions) {
298                    panic!(
299                        "non-zero sst start id {} for NewCompactionGroup",
300                        new_sst_start_id
301                    );
302                } else {
303                    warn!(
304                        %new_sst_start_id,
305                        "non-zero sst start id for NewCompactionGroup"
306                    );
307                }
308            }
309            return;
310        } else if !self.levels.contains_key(&parent_group_id) {
311            unreachable!(
312                "non-existing parent group id {} to init from",
313                parent_group_id
314            );
315        }
316        let [parent_levels, cur_levels] = self
317            .levels
318            .get_disjoint_mut([&parent_group_id, &group_id])
319            .map(|res| res.unwrap());
320        // After certain compaction group operation, e.g. split, any ongoing compaction tasks created prior to that should be rejected due to expiration.
321        // By incrementing the compaction_group_version_id of the compaction group, and comparing it with the one recorded in compaction task, expired compaction tasks can be identified.
322        parent_levels.compaction_group_version_id += 1;
323        cur_levels.compaction_group_version_id += 1;
324        let l0 = &mut parent_levels.l0;
325        {
326            for sub_level in &mut l0.sub_levels {
327                let target_l0 = &mut cur_levels.l0;
328                // Remove SST from sub level may result in empty sub level. It will be purged
329                // whenever another compaction task is finished.
330                let insert_table_infos =
331                    split_sst_info_for_level(&member_table_ids, sub_level, &mut new_sst_id);
332                sub_level.normalize();
333                if insert_table_infos.is_empty() {
334                    continue;
335                }
336                match group_split::get_sub_level_insert_hint(&target_l0.sub_levels, sub_level) {
337                    Ok(idx) => {
338                        add_ssts_to_sub_level(target_l0, idx, insert_table_infos);
339                    }
340                    Err(idx) => {
341                        insert_new_sub_level(
342                            target_l0,
343                            sub_level.sub_level_id,
344                            sub_level.level_type,
345                            insert_table_infos,
346                            Some(idx),
347                        );
348                    }
349                }
350            }
351            l0.normalize();
352        }
353        for (idx, level) in parent_levels.levels.iter_mut().enumerate() {
354            let insert_table_infos =
355                split_sst_info_for_level(&member_table_ids, level, &mut new_sst_id);
356            cur_levels.levels[idx].total_file_size += insert_table_infos
357                .iter()
358                .map(|sst| sst.sst_size)
359                .sum::<u64>();
360            cur_levels.levels[idx].uncompressed_file_size += insert_table_infos
361                .iter()
362                .map(|sst| sst.uncompressed_file_size)
363                .sum::<u64>();
364            cur_levels.levels[idx]
365                .table_infos
366                .extend(insert_table_infos);
367            cur_levels.levels[idx]
368                .table_infos
369                .sort_by(|sst1, sst2| sst1.key_range.cmp(&sst2.key_range));
370            assert!(can_concat(&cur_levels.levels[idx].table_infos));
371            level.normalize();
372        }
373
374        assert!(
375            parent_levels
376                .l0
377                .sub_levels
378                .iter()
379                .all(|level| !level.table_infos.is_empty())
380        );
381        assert!(
382            cur_levels
383                .l0
384                .sub_levels
385                .iter()
386                .all(|level| !level.table_infos.is_empty())
387        );
388    }
389
390    pub fn build_sst_delta_infos(
391        &self,
392        version_delta: &HummockVersionDeltaCommon<SstableInfo>,
393    ) -> Vec<SstDeltaInfo> {
394        let mut infos = vec![];
395
396        // Skip trivial move delta for refiller
397        // The trivial move task only changes the position of the sst in the lsm, it does not modify the object information corresponding to the sst, and does not need to re-execute the refill.
398        if version_delta.trivial_move {
399            return infos;
400        }
401
402        for (group_id, group_deltas) in &version_delta.group_deltas {
403            let mut info = SstDeltaInfo::default();
404
405            let mut removed_l0_ssts: BTreeSet<HummockSstableId> = BTreeSet::new();
406            let mut removed_ssts: BTreeMap<u32, BTreeSet<HummockSstableId>> = BTreeMap::new();
407
408            // Build only if all deltas are intra level deltas.
409            if !group_deltas.group_deltas.iter().all(|delta| {
410                matches!(
411                    delta,
412                    GroupDelta::IntraLevel(_) | GroupDelta::NewL0SubLevel(_)
413                )
414            }) {
415                continue;
416            }
417
418            // TODO(MrCroxx): At most one insert delta is allowed here. It's okay for now with the
419            // current `hummock::manager::gen_version_delta` implementation. Better refactor the
420            // struct to reduce conventions.
421            for group_delta in &group_deltas.group_deltas {
422                match group_delta {
423                    GroupDeltaCommon::NewL0SubLevel(inserted_table_infos) => {
424                        if !inserted_table_infos.is_empty() {
425                            info.insert_sst_level = 0;
426                            info.insert_sst_infos
427                                .extend(inserted_table_infos.iter().cloned());
428                        }
429                    }
430                    GroupDeltaCommon::IntraLevel(intra_level) => {
431                        if !intra_level.inserted_table_infos.is_empty() {
432                            info.insert_sst_level = intra_level.level_idx;
433                            info.insert_sst_infos
434                                .extend(intra_level.inserted_table_infos.iter().cloned());
435                        }
436                        if !intra_level.removed_table_ids.is_empty() {
437                            for id in &intra_level.removed_table_ids {
438                                if intra_level.level_idx == 0 {
439                                    removed_l0_ssts.insert(*id);
440                                } else {
441                                    removed_ssts
442                                        .entry(intra_level.level_idx)
443                                        .or_default()
444                                        .insert(*id);
445                                }
446                            }
447                        }
448                    }
449                    GroupDeltaCommon::GroupConstruct(_)
450                    | GroupDeltaCommon::GroupDestroy(_)
451                    | GroupDeltaCommon::GroupMerge(_)
452                    | GroupDeltaCommon::PruneTableIdsFromSsts(_) => {}
453                }
454            }
455
456            let group = self.levels.get(group_id).unwrap();
457            if !removed_l0_ssts.is_empty() {
458                for l0_sub_level in &group.level0().sub_levels {
459                    for sst_info in &l0_sub_level.table_infos {
460                        if removed_l0_ssts.remove(&sst_info.sst_id) {
461                            info.delete_sst_object_ids.push(sst_info.object_id);
462                        }
463                    }
464                }
465            }
466            for level in &group.levels {
467                if let Some(mut removed_level_ssts) = removed_ssts.remove(&level.level_idx) {
468                    for sst_info in &level.table_infos {
469                        if removed_level_ssts.remove(&sst_info.sst_id) {
470                            info.delete_sst_object_ids.push(sst_info.object_id);
471                        }
472                    }
473                    if !removed_level_ssts.is_empty() {
474                        tracing::error!(
475                            "removed_level_ssts is not empty: {:?}",
476                            removed_level_ssts,
477                        );
478                    }
479                    debug_assert!(removed_level_ssts.is_empty());
480                }
481            }
482
483            if !removed_l0_ssts.is_empty() || !removed_ssts.is_empty() {
484                tracing::error!(
485                    "not empty removed_l0_ssts: {:?}, removed_ssts: {:?}",
486                    removed_l0_ssts,
487                    removed_ssts
488                );
489            }
490            debug_assert!(removed_l0_ssts.is_empty());
491            debug_assert!(removed_ssts.is_empty());
492
493            infos.push(info);
494        }
495
496        infos
497    }
498
499    pub fn apply_version_delta(
500        &mut self,
501        version_delta: &HummockVersionDeltaCommon<SstableInfo>,
502    ) -> HashMap<TableId, Option<StateTableInfo>> {
503        assert_eq!(self.id, version_delta.prev_id);
504
505        let (changed_table_info, mut is_commit_epoch) = self.state_table_info.apply_delta(
506            &version_delta.state_table_info_delta,
507            &version_delta.removed_table_ids,
508        );
509
510        #[expect(deprecated)]
511        {
512            if !is_commit_epoch && self.max_committed_epoch < version_delta.max_committed_epoch {
513                is_commit_epoch = true;
514                tracing::trace!(
515                    "max committed epoch bumped but no table committed epoch is changed"
516                );
517            }
518        }
519
520        // apply to `levels`, which is different compaction groups
521        for (compaction_group_id, group_deltas) in &version_delta.group_deltas {
522            let mut is_applied_l0_compact = false;
523            for group_delta in &group_deltas.group_deltas {
524                match group_delta {
525                    GroupDeltaCommon::GroupConstruct(group_construct) => {
526                        let mut new_levels = build_initial_compaction_group_levels(
527                            *compaction_group_id,
528                            group_construct.get_group_config().unwrap(),
529                        );
530                        let parent_group_id = group_construct.parent_group_id;
531                        new_levels.parent_group_id = parent_group_id;
532                        #[expect(deprecated)]
533                        // for backward-compatibility of previous hummock version delta
534                        new_levels
535                            .member_table_ids
536                            .clone_from(&group_construct.table_ids);
537                        self.levels.insert(*compaction_group_id, new_levels);
538                        if group_construct.version() >= CompatibilityVersion::SplitGroupByTableId {
539                            let split_key = group_construct.split_key.clone().map(Bytes::from);
540                            self.init_with_parent_group_v2(
541                                parent_group_id,
542                                *compaction_group_id,
543                                group_construct.new_sst_start_id,
544                                split_key,
545                            );
546                        } else {
547                            // for backward-compatibility of previous hummock version delta
548                            let member_table_ids = if group_construct.version()
549                                >= CompatibilityVersion::NoMemberTableIds
550                            {
551                                self.state_table_info
552                                    .compaction_group_member_table_ids(*compaction_group_id)
553                                    .iter()
554                                    .copied()
555                                    .collect()
556                            } else {
557                                #[expect(deprecated)]
558                                BTreeSet::from_iter(
559                                    group_construct.table_ids.iter().copied().map(Into::into),
560                                )
561                            };
562                            self.init_with_parent_group(
563                                parent_group_id,
564                                *compaction_group_id,
565                                member_table_ids,
566                                group_construct.new_sst_start_id,
567                            );
568                        }
569                    }
570                    GroupDeltaCommon::GroupMerge(group_merge) => {
571                        tracing::info!(
572                            "group_merge left {:?} right {:?}",
573                            group_merge.left_group_id,
574                            group_merge.right_group_id
575                        );
576                        self.merge_compaction_group(
577                            group_merge.left_group_id,
578                            group_merge.right_group_id,
579                        )
580                    }
581                    GroupDeltaCommon::IntraLevel(level_delta) => {
582                        let levels =
583                            self.levels.get_mut(compaction_group_id).unwrap_or_else(|| {
584                                panic!("compaction group {} does not exist", compaction_group_id)
585                            });
586                        if is_commit_epoch {
587                            assert!(
588                                level_delta.removed_table_ids.is_empty(),
589                                "no sst should be deleted when committing an epoch"
590                            );
591
592                            let IntraLevelDelta {
593                                level_idx,
594                                l0_sub_level_id,
595                                inserted_table_infos,
596                                ..
597                            } = level_delta;
598                            {
599                                assert_eq!(
600                                    *level_idx, 0,
601                                    "we should only add to L0 when we commit an epoch."
602                                );
603                                if !inserted_table_infos.is_empty() {
604                                    insert_new_sub_level(
605                                        &mut levels.l0,
606                                        *l0_sub_level_id,
607                                        PbLevelType::Overlapping,
608                                        inserted_table_infos.clone(),
609                                        None,
610                                    );
611                                }
612                            }
613                        } else {
614                            // The delta is caused by compaction.
615                            levels.apply_compact_ssts(
616                                level_delta,
617                                self.state_table_info
618                                    .compaction_group_member_table_ids(*compaction_group_id),
619                            );
620                            if level_delta.level_idx == 0 {
621                                is_applied_l0_compact = true;
622                            }
623                        }
624                    }
625                    GroupDeltaCommon::NewL0SubLevel(inserted_table_infos) => {
626                        let levels =
627                            self.levels.get_mut(compaction_group_id).unwrap_or_else(|| {
628                                panic!("compaction group {} does not exist", compaction_group_id)
629                            });
630                        assert!(is_commit_epoch);
631
632                        if !inserted_table_infos.is_empty() {
633                            let next_l0_sub_level_id = levels
634                                .l0
635                                .sub_levels
636                                .last()
637                                .map(|level| level.sub_level_id + 1)
638                                .unwrap_or(1);
639
640                            insert_new_sub_level(
641                                &mut levels.l0,
642                                next_l0_sub_level_id,
643                                PbLevelType::Overlapping,
644                                inserted_table_infos.clone(),
645                                None,
646                            );
647                        }
648                    }
649                    GroupDeltaCommon::GroupDestroy(_) => {
650                        self.levels.remove(compaction_group_id);
651                    }
652
653                    GroupDeltaCommon::PruneTableIdsFromSsts(table_ids) => {
654                        self.levels
655                            .get_mut(compaction_group_id)
656                            .unwrap_or_else(|| {
657                                panic!("compaction group {} does not exist", compaction_group_id)
658                            })
659                            .prune_table_ids_from_ssts(table_ids);
660                    }
661                }
662            }
663
664            if is_applied_l0_compact && let Some(levels) = self.levels.get_mut(compaction_group_id)
665            {
666                levels.l0.normalize();
667            }
668        }
669        self.id = version_delta.id;
670        #[expect(deprecated)]
671        {
672            self.max_committed_epoch = version_delta.max_committed_epoch;
673        }
674
675        // apply to table watermark
676
677        // Store the table watermarks that needs to be updated. None means to remove the table watermark of the table id
678        let mut modified_table_watermarks: HashMap<TableId, Option<TableWatermarks>> =
679            HashMap::new();
680
681        // apply to table watermark
682        for (table_id, table_watermarks) in &version_delta.new_table_watermarks {
683            if let Some(current_table_watermarks) = self.table_watermarks.get(table_id) {
684                if version_delta.removed_table_ids.contains(table_id) {
685                    modified_table_watermarks.insert(*table_id, None);
686                } else {
687                    let mut current_table_watermarks = (**current_table_watermarks).clone();
688                    current_table_watermarks.apply_new_table_watermarks(table_watermarks);
689                    modified_table_watermarks.insert(*table_id, Some(current_table_watermarks));
690                }
691            } else {
692                modified_table_watermarks.insert(*table_id, Some(table_watermarks.clone()));
693            }
694        }
695        for table_id in &version_delta.removed_table_ids {
696            modified_table_watermarks.insert(*table_id, None);
697        }
698        for (table_id, table_watermarks) in &self.table_watermarks {
699            let safe_epoch = if let Some(state_table_info) =
700                self.state_table_info.info().get(table_id)
701                && let Some((oldest_epoch, _)) = table_watermarks.watermarks.first()
702                && state_table_info.committed_epoch > *oldest_epoch
703            {
704                // safe epoch has progressed, need further clear.
705                state_table_info.committed_epoch
706            } else {
707                // safe epoch not progressed or the table has been removed. No need to truncate
708                continue;
709            };
710            let table_watermarks = modified_table_watermarks
711                .entry(*table_id)
712                .or_insert_with(|| Some((**table_watermarks).clone()));
713            if let Some(table_watermarks) = table_watermarks {
714                table_watermarks.clear_stale_epoch_watermark(safe_epoch);
715            }
716        }
717        // apply the staging table watermark to hummock version
718        for (table_id, table_watermarks) in modified_table_watermarks {
719            if let Some(table_watermarks) = table_watermarks {
720                self.table_watermarks
721                    .insert(table_id, Arc::new(table_watermarks));
722            } else {
723                self.table_watermarks.remove(&table_id);
724            }
725        }
726        // apply to vector index
727        apply_vector_index_delta(
728            &mut self.vector_indexes,
729            &version_delta.vector_index_delta,
730            &version_delta.removed_table_ids,
731        );
732
733        changed_table_info
734    }
735
736    pub fn apply_change_log_delta<T: Clone>(
737        table_change_log: &mut HashMap<TableId, TableChangeLogCommon<T>>,
738        change_log_delta: &HashMap<TableId, EpochNewChangeLogCommon<T>>,
739    ) {
740        for (table_id, new_change_log) in change_log_delta {
741            match table_change_log.entry(*table_id) {
742                Entry::Occupied(entry) => {
743                    let change_log = entry.into_mut();
744                    change_log.add_change_log(new_change_log.clone());
745                }
746                Entry::Vacant(entry) => {
747                    entry.insert(TableChangeLogCommon::new(once(new_change_log.clone())));
748                }
749            };
750        }
751    }
752
753    /// Returns the log deltas required to truncate the entire change log for a table.
754    pub fn collect_gc_change_log_delta<'a, T: Clone>(
755        current_change_log_table_ids: impl Iterator<Item = &'a TableId>,
756        change_log_delta: &HashMap<TableId, EpochNewChangeLogCommon<T>>,
757        removed_table_ids: &HashSet<TableId>,
758        state_table_info_delta: &HashMap<TableId, StateTableInfoDelta>,
759        changed_table_info: &HashMap<TableId, Option<StateTableInfo>>,
760    ) -> HashSet<TableId> {
761        let mut gc_change_log_delta = HashSet::new();
762        // If a table has no new change log entry (even an empty one), it means we have stopped maintained
763        // the change log for the table, and then we will remove the table change log.
764        // The table change log will also be removed when the table id is removed.
765        for table_id in current_change_log_table_ids {
766            if removed_table_ids.contains(table_id) {
767                gc_change_log_delta.insert(*table_id);
768                continue;
769            }
770            if let Some(table_info_delta) = state_table_info_delta.get(table_id)
771                && let Some(Some(prev_table_info)) = changed_table_info.get(table_id)
772                && table_info_delta.committed_epoch > prev_table_info.committed_epoch
773            {
774                // the table exists previously, and its committed epoch has progressed.
775            } else {
776                // otherwise, the table change log should be kept anyway
777                continue;
778            }
779            let contains = change_log_delta.contains_key(table_id);
780            if !contains {
781                gc_change_log_delta.insert(*table_id);
782                static LOG_SUPPRESSOR: LazyLock<LogSuppressor> =
783                    LazyLock::new(|| LogSuppressor::per_second(1));
784                if let Ok(suppressed_count) = LOG_SUPPRESSOR.check() {
785                    warn!(
786                        suppressed_count,
787                        %table_id,
788                        "table change log dropped due to no further change log at newly committed epoch"
789                    );
790                }
791            }
792        }
793        gc_change_log_delta
794    }
795
796    pub fn build_branched_sst_info(&self) -> BTreeMap<HummockSstableObjectId, BranchedSstInfo> {
797        let mut ret: BTreeMap<_, _> = BTreeMap::new();
798        for (compaction_group_id, group) in &self.levels {
799            for level in group.l0.sub_levels.iter().chain(&group.levels) {
800                for table_info in &level.table_infos {
801                    if table_info.sst_id.as_raw_id() == table_info.object_id.as_raw_id() {
802                        continue;
803                    }
804                    let object_id = table_info.object_id;
805                    let entry: &mut BranchedSstInfo = ret.entry(object_id).or_default();
806                    entry
807                        .entry(*compaction_group_id)
808                        .or_default()
809                        .push(table_info.sst_id)
810                }
811            }
812        }
813        ret
814    }
815
816    pub fn merge_compaction_group(
817        &mut self,
818        left_group_id: CompactionGroupId,
819        right_group_id: CompactionGroupId,
820    ) {
821        // Double check
822        let left_group_id_table_ids = self
823            .state_table_info
824            .compaction_group_member_table_ids(left_group_id)
825            .iter();
826        let right_group_id_table_ids = self
827            .state_table_info
828            .compaction_group_member_table_ids(right_group_id)
829            .iter();
830
831        assert!(
832            left_group_id_table_ids
833                .chain(right_group_id_table_ids)
834                .is_sorted()
835        );
836
837        let total_cg = self.levels.keys().cloned().collect::<Vec<_>>();
838        let right_levels = self.levels.remove(&right_group_id).unwrap_or_else(|| {
839            panic!(
840                "compaction group should exist right {} all {:?}",
841                right_group_id, total_cg
842            )
843        });
844
845        let left_levels = self.levels.get_mut(&left_group_id).unwrap_or_else(|| {
846            panic!(
847                "compaction group should exist left {} all {:?}",
848                left_group_id, total_cg
849            )
850        });
851
852        group_split::merge_levels(left_levels, right_levels);
853    }
854
855    pub fn init_with_parent_group_v2(
856        &mut self,
857        parent_group_id: CompactionGroupId,
858        group_id: CompactionGroupId,
859        new_sst_start_id: HummockSstableId,
860        split_key: Option<Bytes>,
861    ) {
862        let mut new_sst_id = new_sst_start_id;
863        if parent_group_id == StaticCompactionGroupId::NewCompactionGroup {
864            if new_sst_start_id != 0 {
865                if cfg!(debug_assertions) {
866                    panic!(
867                        "non-zero sst start id {} for NewCompactionGroup",
868                        new_sst_start_id
869                    );
870                } else {
871                    warn!(
872                        %new_sst_start_id,
873                        "non-zero sst start id for NewCompactionGroup"
874                    );
875                }
876            }
877            return;
878        } else if !self.levels.contains_key(&parent_group_id) {
879            unreachable!(
880                "non-existing parent group id {} to init from (V2)",
881                parent_group_id
882            );
883        }
884
885        let [parent_levels, cur_levels] = self
886            .levels
887            .get_disjoint_mut([&parent_group_id, &group_id])
888            .map(|res| res.unwrap());
889        // After certain compaction group operation, e.g. split, any ongoing compaction tasks created prior to that should be rejected due to expiration.
890        // By incrementing the compaction_group_version_id of the compaction group, and comparing it with the one recorded in compaction task, expired compaction tasks can be identified.
891        parent_levels.compaction_group_version_id += 1;
892        cur_levels.compaction_group_version_id += 1;
893
894        let l0 = &mut parent_levels.l0;
895        {
896            for sub_level in &mut l0.sub_levels {
897                let target_l0 = &mut cur_levels.l0;
898                // Remove SST from sub level may result in empty sub level. It will be purged
899                // whenever another compaction task is finished.
900                let insert_table_infos = if let Some(split_key) = &split_key {
901                    group_split::split_sst_info_for_level_v2(
902                        sub_level,
903                        &mut new_sst_id,
904                        split_key.clone(),
905                    )
906                } else {
907                    vec![]
908                };
909
910                if insert_table_infos.is_empty() {
911                    continue;
912                }
913
914                sub_level.normalize();
915                match group_split::get_sub_level_insert_hint(&target_l0.sub_levels, sub_level) {
916                    Ok(idx) => {
917                        add_ssts_to_sub_level(target_l0, idx, insert_table_infos);
918                    }
919                    Err(idx) => {
920                        insert_new_sub_level(
921                            target_l0,
922                            sub_level.sub_level_id,
923                            sub_level.level_type,
924                            insert_table_infos,
925                            Some(idx),
926                        );
927                    }
928                }
929            }
930            l0.normalize();
931        }
932
933        for (idx, level) in parent_levels.levels.iter_mut().enumerate() {
934            let insert_table_infos = if let Some(split_key) = &split_key {
935                group_split::split_sst_info_for_level_v2(level, &mut new_sst_id, split_key.clone())
936            } else {
937                vec![]
938            };
939
940            if insert_table_infos.is_empty() {
941                continue;
942            }
943
944            cur_levels.levels[idx].total_file_size += insert_table_infos
945                .iter()
946                .map(|sst| sst.sst_size)
947                .sum::<u64>();
948            cur_levels.levels[idx].uncompressed_file_size += insert_table_infos
949                .iter()
950                .map(|sst| sst.uncompressed_file_size)
951                .sum::<u64>();
952            cur_levels.levels[idx]
953                .table_infos
954                .extend(insert_table_infos);
955            cur_levels.levels[idx]
956                .table_infos
957                .sort_by(|sst1, sst2| sst1.key_range.cmp(&sst2.key_range));
958            assert!(can_concat(&cur_levels.levels[idx].table_infos));
959            level.normalize();
960        }
961
962        assert!(
963            parent_levels
964                .l0
965                .sub_levels
966                .iter()
967                .all(|level| !level.table_infos.is_empty())
968        );
969        assert!(
970            cur_levels
971                .l0
972                .sub_levels
973                .iter()
974                .all(|level| !level.table_infos.is_empty())
975        );
976    }
977}
978
979impl<T> HummockVersionCommon<T>
980where
981    T: SstableIdReader + ObjectIdReader,
982{
983    pub fn get_object_ids(&self) -> impl Iterator<Item = HummockObjectId> + '_ {
984        // DO NOT REMOVE THIS LINE
985        // This is to ensure that when adding new variant to `HummockObjectId`,
986        // the compiler will warn us if we forget to handle it here.
987        match HummockObjectId::Sstable(0.into()) {
988            HummockObjectId::Sstable(_) => {}
989            HummockObjectId::VectorFile(_) => {}
990            HummockObjectId::HnswGraphFile(_) => {}
991        };
992        self.get_sst_infos()
993            .map(|s| HummockObjectId::Sstable(s.object_id()))
994            .chain(
995                self.vector_indexes
996                    .values()
997                    .flat_map(|index| index.get_objects().map(|(object_id, _)| object_id)),
998            )
999    }
1000
1001    pub fn get_sst_ids(&self) -> HashSet<HummockSstableId> {
1002        self.get_sst_infos().map(|s| s.sst_id()).collect()
1003    }
1004
1005    pub fn get_sst_infos(&self) -> impl Iterator<Item = &T> {
1006        self.get_combined_levels()
1007            .flat_map(|level| level.table_infos.iter())
1008    }
1009}
1010
1011impl Levels {
1012    pub(crate) fn apply_compact_ssts(
1013        &mut self,
1014        level_delta: &IntraLevelDeltaCommon<SstableInfo>,
1015        member_table_ids: &BTreeSet<TableId>,
1016    ) {
1017        let IntraLevelDeltaCommon {
1018            level_idx,
1019            l0_sub_level_id,
1020            inserted_table_infos: insert_table_infos,
1021            vnode_partition_count,
1022            removed_table_ids: delete_sst_ids_set,
1023            compaction_group_version_id,
1024        } = level_delta;
1025        let new_vnode_partition_count = *vnode_partition_count;
1026
1027        if is_compaction_task_expired(
1028            self.compaction_group_version_id,
1029            *compaction_group_version_id,
1030        ) {
1031            warn!(
1032                current_compaction_group_version_id = self.compaction_group_version_id,
1033                delta_compaction_group_version_id = compaction_group_version_id,
1034                level_idx,
1035                l0_sub_level_id,
1036                insert_table_infos = ?insert_table_infos
1037                    .iter()
1038                    .map(|sst| (sst.sst_id, sst.object_id))
1039                    .collect_vec(),
1040                ?delete_sst_ids_set,
1041                "This VersionDelta may be committed by an expired compact task. Please check it."
1042            );
1043            return;
1044        }
1045        if !delete_sst_ids_set.is_empty() {
1046            if *level_idx == 0 {
1047                for level in &mut self.l0.sub_levels {
1048                    level.delete_ssts(delete_sst_ids_set);
1049                }
1050            } else {
1051                let idx = *level_idx as usize - 1;
1052                self.levels[idx].delete_ssts(delete_sst_ids_set);
1053            }
1054        }
1055
1056        if !insert_table_infos.is_empty() {
1057            let insert_sst_level_id = *level_idx;
1058            let insert_sub_level_id = *l0_sub_level_id;
1059            if insert_sst_level_id == 0 {
1060                let l0 = &mut self.l0;
1061                let index = l0
1062                    .sub_levels
1063                    .partition_point(|level| level.sub_level_id < insert_sub_level_id);
1064                assert!(
1065                    index < l0.sub_levels.len()
1066                        && l0.sub_levels[index].sub_level_id == insert_sub_level_id,
1067                    "should find the level to insert into when applying compaction generated delta. sub level idx: {},  removed sst ids: {:?}, sub levels: {:?},",
1068                    insert_sub_level_id,
1069                    delete_sst_ids_set,
1070                    l0.sub_levels
1071                        .iter()
1072                        .map(|level| level.sub_level_id)
1073                        .collect_vec()
1074                );
1075                if l0.sub_levels[index].table_infos.is_empty()
1076                    && member_table_ids.len() == 1
1077                    && insert_table_infos.iter().all(|sst| {
1078                        sst.table_ids.len() == 1
1079                            && sst.table_ids[0]
1080                                == *member_table_ids.iter().next().expect("non-empty")
1081                    })
1082                {
1083                    // Only change vnode_partition_count for group which has only one state-table.
1084                    // Only change vnode_partition_count for level which update all sst files in this compact task.
1085                    l0.sub_levels[index].vnode_partition_count = new_vnode_partition_count;
1086                }
1087                level_insert_ssts(&mut l0.sub_levels[index], insert_table_infos);
1088            } else {
1089                let idx = insert_sst_level_id as usize - 1;
1090                if self.levels[idx].table_infos.is_empty()
1091                    && insert_table_infos
1092                        .iter()
1093                        .all(|sst| sst.table_ids.len() == 1)
1094                {
1095                    self.levels[idx].vnode_partition_count = new_vnode_partition_count;
1096                } else if self.levels[idx].vnode_partition_count != 0
1097                    && new_vnode_partition_count == 0
1098                    && member_table_ids.len() > 1
1099                {
1100                    self.levels[idx].vnode_partition_count = 0;
1101                }
1102                level_insert_ssts(&mut self.levels[idx], insert_table_infos);
1103            }
1104        }
1105    }
1106
1107    /// Prune specified table ids from all SST metadata, remove emptied SSTs and sub-levels,
1108    /// then bump `compaction_group_version_id`.
1109    pub(crate) fn prune_table_ids_from_ssts(&mut self, table_ids: &HashSet<TableId>) {
1110        for level in self.l0.sub_levels.iter_mut().chain(self.levels.iter_mut()) {
1111            level.prune_table_ids_from_ssts(table_ids);
1112        }
1113        self.l0.normalize();
1114        self.compaction_group_version_id += 1;
1115    }
1116}
1117
1118impl<T> HummockVersionCommon<T> {
1119    pub fn get_combined_levels(&self) -> impl Iterator<Item = &'_ LevelCommon<T>> + '_ {
1120        self.levels
1121            .values()
1122            .flat_map(|level| level.l0.sub_levels.iter().rev().chain(level.levels.iter()))
1123    }
1124}
1125
1126pub fn build_initial_compaction_group_levels(
1127    group_id: impl Into<CompactionGroupId>,
1128    compaction_config: &CompactionConfig,
1129) -> Levels {
1130    let mut levels = vec![];
1131    for l in 0..compaction_config.get_max_level() {
1132        levels.push(Level {
1133            level_idx: (l + 1) as u32,
1134            level_type: PbLevelType::Nonoverlapping,
1135            table_infos: vec![],
1136            total_file_size: 0,
1137            sub_level_id: 0,
1138            uncompressed_file_size: 0,
1139            vnode_partition_count: 0,
1140        });
1141    }
1142    #[expect(deprecated)] // for backward-compatibility of previous hummock version delta
1143    Levels {
1144        levels,
1145        l0: OverlappingLevel {
1146            sub_levels: vec![],
1147            total_file_size: 0,
1148            uncompressed_file_size: 0,
1149        },
1150        group_id: group_id.into(),
1151        parent_group_id: 0.into(),
1152        member_table_ids: vec![],
1153        compaction_group_version_id: 0,
1154    }
1155}
1156
1157fn split_sst_info_for_level(
1158    member_table_ids: &BTreeSet<TableId>,
1159    level: &mut Level,
1160    new_sst_id: &mut HummockSstableId,
1161) -> Vec<SstableInfo> {
1162    // Remove SST from sub level may result in empty sub level. It will be purged
1163    // whenever another compaction task is finished.
1164    let mut insert_table_infos = vec![];
1165    for sst_info in &mut level.table_infos {
1166        let contains_member_table = sst_info
1167            .table_ids
1168            .iter()
1169            .any(|table_id| member_table_ids.contains(table_id));
1170        let sst_size = sst_info.sst_size;
1171        if sst_size / 2 == 0 {
1172            tracing::warn!(
1173                id = %sst_info.sst_id,
1174                object_id = %sst_info.object_id,
1175                sst_size = sst_info.sst_size,
1176                file_size = sst_info.file_size,
1177                "Sstable sst_size is under expected",
1178            );
1179        };
1180        if contains_member_table {
1181            let (modified_sst, branch_sst) = split_sst_with_table_ids(
1182                sst_info,
1183                new_sst_id,
1184                sst_size / 2,
1185                sst_size / 2,
1186                member_table_ids.iter().cloned().collect_vec(),
1187            );
1188            *sst_info = modified_sst;
1189            insert_table_infos.push(branch_sst);
1190        }
1191    }
1192    insert_table_infos
1193}
1194
1195/// Gets all compaction group ids.
1196pub fn get_compaction_group_ids(
1197    version: &HummockVersion,
1198) -> impl Iterator<Item = CompactionGroupId> + '_ {
1199    version.levels.keys().cloned()
1200}
1201
1202pub fn get_table_compaction_group_id_mapping(
1203    version: &HummockVersion,
1204) -> HashMap<StateTableId, CompactionGroupId> {
1205    version
1206        .state_table_info
1207        .info()
1208        .iter()
1209        .map(|(table_id, info)| (*table_id, info.compaction_group_id))
1210        .collect()
1211}
1212
1213/// Gets all SSTs in `group_id`
1214pub fn get_compaction_group_ssts(
1215    version: &HummockVersion,
1216    group_id: CompactionGroupId,
1217) -> impl Iterator<Item = (HummockSstableObjectId, HummockSstableId)> + '_ {
1218    let group_levels = version.get_compaction_group_levels(group_id);
1219    group_levels
1220        .l0
1221        .sub_levels
1222        .iter()
1223        .rev()
1224        .chain(group_levels.levels.iter())
1225        .flat_map(|level| {
1226            level
1227                .table_infos
1228                .iter()
1229                .map(|table_info| (table_info.object_id, table_info.sst_id))
1230        })
1231}
1232
1233pub fn new_sub_level(
1234    sub_level_id: u64,
1235    level_type: PbLevelType,
1236    table_infos: Vec<SstableInfo>,
1237) -> Level {
1238    if level_type == PbLevelType::Nonoverlapping {
1239        debug_assert!(
1240            can_concat(&table_infos),
1241            "sst of non-overlapping level is not concat-able: {:?}",
1242            table_infos
1243        );
1244    }
1245    let total_file_size = table_infos.iter().map(|table| table.sst_size).sum();
1246    let uncompressed_file_size = table_infos
1247        .iter()
1248        .map(|table| table.uncompressed_file_size)
1249        .sum();
1250    Level {
1251        level_idx: 0,
1252        level_type,
1253        table_infos,
1254        total_file_size,
1255        sub_level_id,
1256        uncompressed_file_size,
1257        vnode_partition_count: 0,
1258    }
1259}
1260
1261pub fn add_ssts_to_sub_level(
1262    l0: &mut OverlappingLevel,
1263    sub_level_idx: usize,
1264    insert_table_infos: Vec<SstableInfo>,
1265) {
1266    insert_table_infos.iter().for_each(|sst| {
1267        l0.sub_levels[sub_level_idx].total_file_size += sst.sst_size;
1268        l0.sub_levels[sub_level_idx].uncompressed_file_size += sst.uncompressed_file_size;
1269        l0.total_file_size += sst.sst_size;
1270        l0.uncompressed_file_size += sst.uncompressed_file_size;
1271    });
1272    l0.sub_levels[sub_level_idx]
1273        .table_infos
1274        .extend(insert_table_infos);
1275    if l0.sub_levels[sub_level_idx].level_type == PbLevelType::Nonoverlapping {
1276        l0.sub_levels[sub_level_idx]
1277            .table_infos
1278            .sort_by(|sst1, sst2| sst1.key_range.cmp(&sst2.key_range));
1279        assert!(
1280            can_concat(&l0.sub_levels[sub_level_idx].table_infos),
1281            "sstable ids: {:?}",
1282            l0.sub_levels[sub_level_idx]
1283                .table_infos
1284                .iter()
1285                .map(|sst| sst.sst_id)
1286                .collect_vec()
1287        );
1288    }
1289}
1290
1291/// `None` value of `sub_level_insert_hint` means append.
1292pub fn insert_new_sub_level(
1293    l0: &mut OverlappingLevel,
1294    insert_sub_level_id: u64,
1295    level_type: PbLevelType,
1296    insert_table_infos: Vec<SstableInfo>,
1297    sub_level_insert_hint: Option<usize>,
1298) {
1299    if insert_sub_level_id == u64::MAX {
1300        return;
1301    }
1302    let insert_pos = if let Some(insert_pos) = sub_level_insert_hint {
1303        insert_pos
1304    } else {
1305        if let Some(newest_level) = l0.sub_levels.last() {
1306            assert!(
1307                newest_level.sub_level_id < insert_sub_level_id,
1308                "inserted new level is not the newest: prev newest: {}, insert: {}. L0: {:?}",
1309                newest_level.sub_level_id,
1310                insert_sub_level_id,
1311                l0,
1312            );
1313        }
1314        l0.sub_levels.len()
1315    };
1316    #[cfg(debug_assertions)]
1317    {
1318        if insert_pos > 0
1319            && let Some(smaller_level) = l0.sub_levels.get(insert_pos - 1)
1320        {
1321            debug_assert!(smaller_level.sub_level_id < insert_sub_level_id);
1322        }
1323        if let Some(larger_level) = l0.sub_levels.get(insert_pos) {
1324            debug_assert!(larger_level.sub_level_id > insert_sub_level_id);
1325        }
1326    }
1327    // All files will be committed in one new Overlapping sub-level and become
1328    // Nonoverlapping  after at least one compaction.
1329    let level = new_sub_level(insert_sub_level_id, level_type, insert_table_infos);
1330    l0.total_file_size += level.total_file_size;
1331    l0.uncompressed_file_size += level.uncompressed_file_size;
1332    l0.sub_levels.insert(insert_pos, level);
1333}
1334
1335impl Level {
1336    fn recompute_size(&mut self) {
1337        self.total_file_size = self
1338            .table_infos
1339            .iter()
1340            .map(|table| table.sst_size)
1341            .sum::<u64>();
1342        self.uncompressed_file_size = self
1343            .table_infos
1344            .iter()
1345            .map(|table| table.uncompressed_file_size)
1346            .sum::<u64>();
1347    }
1348
1349    /// Remove SSTs with empty `table_ids`, then recompute sizes.
1350    fn normalize(&mut self) {
1351        self.table_infos
1352            .retain(|sst_info| !sst_info.table_ids.is_empty());
1353        self.recompute_size();
1354    }
1355
1356    /// Return `true` if any SST was actually removed.
1357    fn delete_ssts(&mut self, ids: &HashSet<HummockSstableId>) -> bool {
1358        let original_len = self.table_infos.len();
1359        self.table_infos
1360            .retain(|table| !ids.contains(&table.sst_id));
1361        self.recompute_size();
1362        original_len != self.table_infos.len()
1363    }
1364
1365    /// Prune specified `table_ids` from each SST's metadata,
1366    /// remove SSTs that become empty, then recompute sizes.
1367    fn prune_table_ids_from_ssts(&mut self, table_ids: &HashSet<TableId>) {
1368        for sstable_info in &mut self.table_infos {
1369            if !sstable_info
1370                .table_ids
1371                .iter()
1372                .any(|table_id| table_ids.contains(table_id))
1373            {
1374                continue;
1375            }
1376
1377            let mut inner = sstable_info.get_inner();
1378            inner.table_ids.retain(|id| !table_ids.contains(id));
1379            sstable_info.set_inner(inner);
1380        }
1381        self.normalize();
1382    }
1383}
1384
1385impl OverlappingLevel {
1386    /// Remove empty sub-levels, then recompute aggregated sizes.
1387    fn normalize(&mut self) {
1388        self.sub_levels
1389            .retain(|level| !level.table_infos.is_empty());
1390        self.total_file_size = self
1391            .sub_levels
1392            .iter()
1393            .map(|level| level.total_file_size)
1394            .sum::<u64>();
1395        self.uncompressed_file_size = self
1396            .sub_levels
1397            .iter()
1398            .map(|level| level.uncompressed_file_size)
1399            .sum::<u64>();
1400    }
1401}
1402
1403fn level_insert_ssts(operand: &mut Level, insert_table_infos: &Vec<SstableInfo>) {
1404    fn display_sstable_infos(ssts: &[impl Borrow<SstableInfo>]) -> String {
1405        format!(
1406            "sstable ids: {:?}",
1407            ssts.iter().map(|s| s.borrow().sst_id).collect_vec()
1408        )
1409    }
1410    operand.total_file_size += insert_table_infos
1411        .iter()
1412        .map(|sst| sst.sst_size)
1413        .sum::<u64>();
1414    operand.uncompressed_file_size += insert_table_infos
1415        .iter()
1416        .map(|sst| sst.uncompressed_file_size)
1417        .sum::<u64>();
1418    if operand.level_type == PbLevelType::Overlapping {
1419        operand.level_type = PbLevelType::Nonoverlapping;
1420        operand
1421            .table_infos
1422            .extend(insert_table_infos.iter().cloned());
1423        operand
1424            .table_infos
1425            .sort_by(|sst1, sst2| sst1.key_range.cmp(&sst2.key_range));
1426        assert!(
1427            can_concat(&operand.table_infos),
1428            "{}",
1429            display_sstable_infos(&operand.table_infos)
1430        );
1431    } else if !insert_table_infos.is_empty() {
1432        let sorted_insert: Vec<_> = insert_table_infos
1433            .iter()
1434            .sorted_by(|sst1, sst2| sst1.key_range.cmp(&sst2.key_range))
1435            .cloned()
1436            .collect();
1437        let first = &sorted_insert[0];
1438        let last = &sorted_insert[sorted_insert.len() - 1];
1439        let pos = operand
1440            .table_infos
1441            .partition_point(|b| b.key_range.cmp(&first.key_range) == Ordering::Less);
1442        if pos >= operand.table_infos.len()
1443            || last.key_range.cmp(&operand.table_infos[pos].key_range) == Ordering::Less
1444        {
1445            operand.table_infos.splice(pos..pos, sorted_insert);
1446            // Validate the inserted SST batch along with the two SSTs that precede and follow it.
1447            let start = pos.saturating_sub(1);
1448            let end = (start + insert_table_infos.len() + 2).min(operand.table_infos.len());
1449            let validate_range = &operand.table_infos[start..end];
1450            assert!(
1451                can_concat(validate_range),
1452                "{}",
1453                display_sstable_infos(validate_range),
1454            );
1455        } else {
1456            // If this branch is reached, it indicates some unexpected behavior in compaction.
1457            // Here we issue a warning and fall back to insert one by one.
1458            warn!(insert = ?insert_table_infos, level = ?operand.table_infos, "unexpected overlap");
1459            for i in insert_table_infos {
1460                let pos = operand
1461                    .table_infos
1462                    .partition_point(|b| b.key_range.cmp(&i.key_range) == Ordering::Less);
1463                operand.table_infos.insert(pos, i.clone());
1464            }
1465            assert!(
1466                can_concat(&operand.table_infos),
1467                "{}",
1468                display_sstable_infos(&operand.table_infos)
1469            );
1470        }
1471    }
1472}
1473
1474pub fn version_object_size_map(version: &HummockVersion) -> HashMap<HummockObjectId, u64> {
1475    // DO NOT REMOVE THIS LINE
1476    // This is to ensure that when adding new variant to `HummockObjectId`,
1477    // the compiler will warn us if we forget to handle it here.
1478    match HummockObjectId::Sstable(0.into()) {
1479        HummockObjectId::Sstable(_) => {}
1480        HummockObjectId::VectorFile(_) => {}
1481        HummockObjectId::HnswGraphFile(_) => {}
1482    };
1483    version
1484        .levels
1485        .values()
1486        .flat_map(|cg| {
1487            cg.level0()
1488                .sub_levels
1489                .iter()
1490                .chain(cg.levels.iter())
1491                .flat_map(|level| level.table_infos.iter().map(|t| (t.object_id, t.file_size)))
1492        })
1493        .map(|(object_id, size)| (HummockObjectId::Sstable(object_id), size))
1494        .chain(
1495            version
1496                .vector_indexes
1497                .values()
1498                .flat_map(|index| index.get_objects()),
1499        )
1500        .collect()
1501}
1502
1503/// Verify the validity of a `HummockVersion` and return a list of violations if any.
1504/// Currently this method is only used by risectl validate-version.
1505pub fn validate_version(version: &HummockVersion) -> Vec<String> {
1506    let mut res = Vec::new();
1507    // Ensure each table maps to only one compaction group
1508    for (group_id, levels) in &version.levels {
1509        // Ensure compaction group id matches
1510        if levels.group_id != *group_id {
1511            res.push(format!(
1512                "GROUP {}: inconsistent group id {} in Levels",
1513                group_id, levels.group_id
1514            ));
1515        }
1516
1517        let validate_level = |group: CompactionGroupId,
1518                              expected_level_idx: u32,
1519                              level: &Level,
1520                              res: &mut Vec<String>| {
1521            let mut level_identifier = format!("GROUP {} LEVEL {}", group, level.level_idx);
1522            if level.level_idx == 0 {
1523                level_identifier.push_str(format!("SUBLEVEL {}", level.sub_level_id).as_str());
1524                // Ensure sub-level is not empty
1525                if level.table_infos.is_empty() {
1526                    res.push(format!("{}: empty level", level_identifier));
1527                }
1528            } else if level.level_type != PbLevelType::Nonoverlapping {
1529                // Ensure non-L0 level is non-overlapping level
1530                res.push(format!(
1531                    "{}: level type {:?} is not non-overlapping",
1532                    level_identifier, level.level_type
1533                ));
1534            }
1535
1536            // Ensure level idx matches
1537            if level.level_idx != expected_level_idx {
1538                res.push(format!(
1539                    "{}: mismatched level idx {}",
1540                    level_identifier, expected_level_idx
1541                ));
1542            }
1543
1544            let mut prev_table_info: Option<&SstableInfo> = None;
1545            for table_info in &level.table_infos {
1546                // Ensure table_ids are sorted and unique
1547                if !table_info.table_ids.is_sorted_by(|a, b| a < b) {
1548                    res.push(format!(
1549                        "{} SST {}: table_ids not sorted",
1550                        level_identifier, table_info.object_id
1551                    ));
1552                }
1553
1554                // Ensure SSTs in non-overlapping level have non-overlapping key range
1555                if level.level_type == PbLevelType::Nonoverlapping {
1556                    if let Some(prev) = prev_table_info.take()
1557                        && prev
1558                            .key_range
1559                            .compare_right_with(&table_info.key_range.left)
1560                            != Ordering::Less
1561                    {
1562                        res.push(format!(
1563                            "{} SST {}: key range should not overlap. prev={:?}, cur={:?}",
1564                            level_identifier, table_info.object_id, prev, table_info
1565                        ));
1566                    }
1567                    let _ = prev_table_info.insert(table_info);
1568                }
1569            }
1570        };
1571
1572        let l0 = &levels.l0;
1573        let mut prev_sub_level_id = None;
1574        for sub_level in &l0.sub_levels {
1575            // Ensure sub_level_id is strictly increasing.
1576            if let Some(prev_sub_level_id) = prev_sub_level_id
1577                && sub_level.sub_level_id <= prev_sub_level_id
1578            {
1579                res.push(format!(
1580                    "GROUP {} LEVEL 0: sub_level_id {} <= prev_sub_level {}",
1581                    group_id, sub_level.sub_level_id, prev_sub_level_id
1582                ));
1583            }
1584            prev_sub_level_id = Some(sub_level.sub_level_id);
1585
1586            validate_level(*group_id, 0, sub_level, &mut res);
1587        }
1588
1589        for idx in 1..=levels.levels.len() {
1590            validate_level(*group_id, idx as u32, levels.get_level(idx), &mut res);
1591        }
1592    }
1593    res
1594}
1595
1596#[cfg(test)]
1597mod tests {
1598    use std::collections::{HashMap, HashSet};
1599    use std::sync::Arc;
1600
1601    use bytes::Bytes;
1602    use risingwave_common::bitmap::Bitmap;
1603    use risingwave_common::catalog::TableId;
1604    use risingwave_common::hash::VirtualNode;
1605    use risingwave_common::util::epoch::test_epoch;
1606    use risingwave_pb::hummock::{
1607        CompactionConfig, GroupConstruct, GroupDestroy, LevelType, StateTableInfo,
1608    };
1609
1610    use super::group_split;
1611    use crate::HummockVersionId;
1612    use crate::compaction_group::group_split::*;
1613    use crate::compaction_group::hummock_version_ext::build_initial_compaction_group_levels;
1614    use crate::key::{FullKey, gen_key_from_str};
1615    use crate::key_range::KeyRange;
1616    use crate::level::{Level, Levels, OverlappingLevel};
1617    use crate::sstable_info::{SstableInfo, SstableInfoInner};
1618    use crate::table_watermark::{
1619        TableWatermarks, VnodeWatermark, WatermarkDirection, WatermarkSerdeType,
1620    };
1621    use crate::version::{
1622        GroupDelta, GroupDeltas, HummockVersion, HummockVersionDelta, HummockVersionStateTableInfo,
1623        IntraLevelDelta,
1624    };
1625
1626    fn gen_sstable_info(sst_id: u64, table_ids: Vec<u32>, epoch: u64) -> SstableInfo {
1627        gen_sstable_info_impl(sst_id, table_ids, epoch).into()
1628    }
1629
1630    fn gen_sstable_info_impl(sst_id: u64, table_ids: Vec<u32>, epoch: u64) -> SstableInfoInner {
1631        let table_key_l = gen_key_from_str(VirtualNode::ZERO, "1");
1632        let table_key_r = gen_key_from_str(VirtualNode::MAX_FOR_TEST, "1");
1633        let full_key_l = FullKey::for_test(
1634            TableId::new(*table_ids.first().unwrap()),
1635            table_key_l,
1636            epoch,
1637        )
1638        .encode();
1639        let full_key_r =
1640            FullKey::for_test(TableId::new(*table_ids.last().unwrap()), table_key_r, epoch)
1641                .encode();
1642
1643        SstableInfoInner {
1644            sst_id: sst_id.into(),
1645            key_range: KeyRange {
1646                left: full_key_l.into(),
1647                right: full_key_r.into(),
1648                right_exclusive: false,
1649            },
1650            table_ids: table_ids.into_iter().map(Into::into).collect(),
1651            object_id: sst_id.into(),
1652            min_epoch: 20,
1653            max_epoch: 20,
1654            file_size: 100,
1655            sst_size: 100,
1656            ..Default::default()
1657        }
1658    }
1659
1660    #[test]
1661    fn test_get_sst_object_ids() {
1662        let mut version = HummockVersion {
1663            id: HummockVersionId::new(0),
1664            levels: HashMap::from_iter([(
1665                0.into(),
1666                Levels {
1667                    levels: vec![],
1668                    l0: OverlappingLevel {
1669                        sub_levels: vec![],
1670                        total_file_size: 0,
1671                        uncompressed_file_size: 0,
1672                    },
1673                    ..Default::default()
1674                },
1675            )]),
1676            ..Default::default()
1677        };
1678        assert_eq!(version.get_object_ids().count(), 0);
1679
1680        // Add to sub level
1681        version
1682            .levels
1683            .get_mut(&0)
1684            .unwrap()
1685            .l0
1686            .sub_levels
1687            .push(Level {
1688                table_infos: vec![
1689                    SstableInfoInner {
1690                        object_id: 11.into(),
1691                        sst_id: 11.into(),
1692                        ..Default::default()
1693                    }
1694                    .into(),
1695                ],
1696                ..Default::default()
1697            });
1698        assert_eq!(version.get_object_ids().count(), 1);
1699
1700        // Add to non sub level
1701        version.levels.get_mut(&0).unwrap().levels.push(Level {
1702            table_infos: vec![
1703                SstableInfoInner {
1704                    object_id: 22.into(),
1705                    sst_id: 22.into(),
1706                    ..Default::default()
1707                }
1708                .into(),
1709            ],
1710            ..Default::default()
1711        });
1712        assert_eq!(version.get_object_ids().count(), 2);
1713    }
1714
1715    #[test]
1716    fn test_apply_version_delta() {
1717        let mut version = HummockVersion {
1718            id: HummockVersionId::new(0),
1719            levels: HashMap::from_iter([
1720                (
1721                    0.into(),
1722                    build_initial_compaction_group_levels(
1723                        0,
1724                        &CompactionConfig {
1725                            max_level: 6,
1726                            ..Default::default()
1727                        },
1728                    ),
1729                ),
1730                (
1731                    1.into(),
1732                    build_initial_compaction_group_levels(
1733                        1,
1734                        &CompactionConfig {
1735                            max_level: 6,
1736                            ..Default::default()
1737                        },
1738                    ),
1739                ),
1740            ]),
1741            ..Default::default()
1742        };
1743        let version_delta = HummockVersionDelta {
1744            id: HummockVersionId::new(1),
1745            group_deltas: HashMap::from_iter([
1746                (
1747                    2.into(),
1748                    GroupDeltas {
1749                        group_deltas: vec![GroupDelta::GroupConstruct(Box::new(GroupConstruct {
1750                            group_config: Some(CompactionConfig {
1751                                max_level: 6,
1752                                ..Default::default()
1753                            }),
1754                            ..Default::default()
1755                        }))],
1756                    },
1757                ),
1758                (
1759                    0.into(),
1760                    GroupDeltas {
1761                        group_deltas: vec![GroupDelta::GroupDestroy(GroupDestroy {})],
1762                    },
1763                ),
1764                (
1765                    1.into(),
1766                    GroupDeltas {
1767                        group_deltas: vec![GroupDelta::IntraLevel(IntraLevelDelta::new(
1768                            1,
1769                            0,
1770                            HashSet::new(),
1771                            vec![
1772                                SstableInfoInner {
1773                                    object_id: 1.into(),
1774                                    sst_id: 1.into(),
1775                                    ..Default::default()
1776                                }
1777                                .into(),
1778                            ],
1779                            0,
1780                            version
1781                                .levels
1782                                .get(&1)
1783                                .as_ref()
1784                                .unwrap()
1785                                .compaction_group_version_id,
1786                        ))],
1787                    },
1788                ),
1789            ]),
1790            ..Default::default()
1791        };
1792        let version_delta = version_delta;
1793
1794        version.apply_version_delta(&version_delta);
1795        let mut cg1 = build_initial_compaction_group_levels(
1796            1,
1797            &CompactionConfig {
1798                max_level: 6,
1799                ..Default::default()
1800            },
1801        );
1802        cg1.levels[0] = Level {
1803            level_idx: 1,
1804            level_type: LevelType::Nonoverlapping,
1805            table_infos: vec![
1806                SstableInfoInner {
1807                    object_id: 1.into(),
1808                    sst_id: 1.into(),
1809                    ..Default::default()
1810                }
1811                .into(),
1812            ],
1813            ..Default::default()
1814        };
1815        assert_eq!(
1816            version,
1817            HummockVersion {
1818                id: HummockVersionId::new(1),
1819                levels: HashMap::from_iter([
1820                    (
1821                        2.into(),
1822                        build_initial_compaction_group_levels(
1823                            2,
1824                            &CompactionConfig {
1825                                max_level: 6,
1826                                ..Default::default()
1827                            },
1828                        ),
1829                    ),
1830                    (1.into(), cg1),
1831                ]),
1832                ..Default::default()
1833            }
1834        );
1835    }
1836
1837    fn gen_sst_info(object_id: u64, table_ids: Vec<u32>, left: Bytes, right: Bytes) -> SstableInfo {
1838        gen_sst_info_impl(object_id, table_ids, left, right).into()
1839    }
1840
1841    fn gen_sst_info_impl(
1842        object_id: u64,
1843        table_ids: Vec<u32>,
1844        left: Bytes,
1845        right: Bytes,
1846    ) -> SstableInfoInner {
1847        SstableInfoInner {
1848            object_id: object_id.into(),
1849            sst_id: object_id.into(),
1850            key_range: KeyRange {
1851                left,
1852                right,
1853                right_exclusive: false,
1854            },
1855            table_ids: table_ids.into_iter().map(Into::into).collect(),
1856            file_size: 100,
1857            sst_size: 100,
1858            uncompressed_file_size: 100,
1859            ..Default::default()
1860        }
1861    }
1862
1863    #[test]
1864    fn test_merge_levels() {
1865        let mut left_levels = build_initial_compaction_group_levels(
1866            1,
1867            &CompactionConfig {
1868                max_level: 6,
1869                ..Default::default()
1870            },
1871        );
1872
1873        let mut right_levels = build_initial_compaction_group_levels(
1874            2,
1875            &CompactionConfig {
1876                max_level: 6,
1877                ..Default::default()
1878            },
1879        );
1880
1881        left_levels.levels[0] = Level {
1882            level_idx: 1,
1883            level_type: LevelType::Nonoverlapping,
1884            table_infos: vec![
1885                gen_sst_info(
1886                    1,
1887                    vec![3],
1888                    FullKey::for_test(
1889                        TableId::new(3),
1890                        gen_key_from_str(VirtualNode::from_index(1), "1"),
1891                        0,
1892                    )
1893                    .encode()
1894                    .into(),
1895                    FullKey::for_test(
1896                        TableId::new(3),
1897                        gen_key_from_str(VirtualNode::from_index(200), "1"),
1898                        0,
1899                    )
1900                    .encode()
1901                    .into(),
1902                ),
1903                gen_sst_info(
1904                    10,
1905                    vec![3, 4],
1906                    FullKey::for_test(
1907                        TableId::new(3),
1908                        gen_key_from_str(VirtualNode::from_index(201), "1"),
1909                        0,
1910                    )
1911                    .encode()
1912                    .into(),
1913                    FullKey::for_test(
1914                        TableId::new(4),
1915                        gen_key_from_str(VirtualNode::from_index(10), "1"),
1916                        0,
1917                    )
1918                    .encode()
1919                    .into(),
1920                ),
1921                gen_sst_info(
1922                    11,
1923                    vec![4],
1924                    FullKey::for_test(
1925                        TableId::new(4),
1926                        gen_key_from_str(VirtualNode::from_index(11), "1"),
1927                        0,
1928                    )
1929                    .encode()
1930                    .into(),
1931                    FullKey::for_test(
1932                        TableId::new(4),
1933                        gen_key_from_str(VirtualNode::from_index(200), "1"),
1934                        0,
1935                    )
1936                    .encode()
1937                    .into(),
1938                ),
1939            ],
1940            total_file_size: 300,
1941            ..Default::default()
1942        };
1943
1944        left_levels.l0.sub_levels.push(Level {
1945            level_idx: 0,
1946            table_infos: vec![gen_sst_info(
1947                3,
1948                vec![3],
1949                FullKey::for_test(
1950                    TableId::new(3),
1951                    gen_key_from_str(VirtualNode::from_index(1), "1"),
1952                    0,
1953                )
1954                .encode()
1955                .into(),
1956                FullKey::for_test(
1957                    TableId::new(3),
1958                    gen_key_from_str(VirtualNode::from_index(200), "1"),
1959                    0,
1960                )
1961                .encode()
1962                .into(),
1963            )],
1964            sub_level_id: 101,
1965            level_type: LevelType::Overlapping,
1966            total_file_size: 100,
1967            ..Default::default()
1968        });
1969
1970        left_levels.l0.sub_levels.push(Level {
1971            level_idx: 0,
1972            table_infos: vec![gen_sst_info(
1973                3,
1974                vec![3],
1975                FullKey::for_test(
1976                    TableId::new(3),
1977                    gen_key_from_str(VirtualNode::from_index(1), "1"),
1978                    0,
1979                )
1980                .encode()
1981                .into(),
1982                FullKey::for_test(
1983                    TableId::new(3),
1984                    gen_key_from_str(VirtualNode::from_index(200), "1"),
1985                    0,
1986                )
1987                .encode()
1988                .into(),
1989            )],
1990            sub_level_id: 103,
1991            level_type: LevelType::Overlapping,
1992            total_file_size: 100,
1993            ..Default::default()
1994        });
1995
1996        left_levels.l0.sub_levels.push(Level {
1997            level_idx: 0,
1998            table_infos: vec![gen_sst_info(
1999                3,
2000                vec![3],
2001                FullKey::for_test(
2002                    TableId::new(3),
2003                    gen_key_from_str(VirtualNode::from_index(1), "1"),
2004                    0,
2005                )
2006                .encode()
2007                .into(),
2008                FullKey::for_test(
2009                    TableId::new(3),
2010                    gen_key_from_str(VirtualNode::from_index(200), "1"),
2011                    0,
2012                )
2013                .encode()
2014                .into(),
2015            )],
2016            sub_level_id: 105,
2017            level_type: LevelType::Nonoverlapping,
2018            total_file_size: 100,
2019            ..Default::default()
2020        });
2021
2022        right_levels.levels[0] = Level {
2023            level_idx: 1,
2024            level_type: LevelType::Nonoverlapping,
2025            table_infos: vec![
2026                gen_sst_info(
2027                    1,
2028                    vec![5],
2029                    FullKey::for_test(
2030                        TableId::new(5),
2031                        gen_key_from_str(VirtualNode::from_index(1), "1"),
2032                        0,
2033                    )
2034                    .encode()
2035                    .into(),
2036                    FullKey::for_test(
2037                        TableId::new(5),
2038                        gen_key_from_str(VirtualNode::from_index(200), "1"),
2039                        0,
2040                    )
2041                    .encode()
2042                    .into(),
2043                ),
2044                gen_sst_info(
2045                    10,
2046                    vec![5, 6],
2047                    FullKey::for_test(
2048                        TableId::new(5),
2049                        gen_key_from_str(VirtualNode::from_index(201), "1"),
2050                        0,
2051                    )
2052                    .encode()
2053                    .into(),
2054                    FullKey::for_test(
2055                        TableId::new(6),
2056                        gen_key_from_str(VirtualNode::from_index(10), "1"),
2057                        0,
2058                    )
2059                    .encode()
2060                    .into(),
2061                ),
2062                gen_sst_info(
2063                    11,
2064                    vec![6],
2065                    FullKey::for_test(
2066                        TableId::new(6),
2067                        gen_key_from_str(VirtualNode::from_index(11), "1"),
2068                        0,
2069                    )
2070                    .encode()
2071                    .into(),
2072                    FullKey::for_test(
2073                        TableId::new(6),
2074                        gen_key_from_str(VirtualNode::from_index(200), "1"),
2075                        0,
2076                    )
2077                    .encode()
2078                    .into(),
2079                ),
2080            ],
2081            total_file_size: 300,
2082            ..Default::default()
2083        };
2084
2085        right_levels.l0.sub_levels.push(Level {
2086            level_idx: 0,
2087            table_infos: vec![gen_sst_info(
2088                3,
2089                vec![5],
2090                FullKey::for_test(
2091                    TableId::new(5),
2092                    gen_key_from_str(VirtualNode::from_index(1), "1"),
2093                    0,
2094                )
2095                .encode()
2096                .into(),
2097                FullKey::for_test(
2098                    TableId::new(5),
2099                    gen_key_from_str(VirtualNode::from_index(200), "1"),
2100                    0,
2101                )
2102                .encode()
2103                .into(),
2104            )],
2105            sub_level_id: 101,
2106            level_type: LevelType::Overlapping,
2107            total_file_size: 100,
2108            ..Default::default()
2109        });
2110
2111        right_levels.l0.sub_levels.push(Level {
2112            level_idx: 0,
2113            table_infos: vec![gen_sst_info(
2114                5,
2115                vec![5],
2116                FullKey::for_test(
2117                    TableId::new(5),
2118                    gen_key_from_str(VirtualNode::from_index(1), "1"),
2119                    0,
2120                )
2121                .encode()
2122                .into(),
2123                FullKey::for_test(
2124                    TableId::new(5),
2125                    gen_key_from_str(VirtualNode::from_index(200), "1"),
2126                    0,
2127                )
2128                .encode()
2129                .into(),
2130            )],
2131            sub_level_id: 102,
2132            level_type: LevelType::Overlapping,
2133            total_file_size: 100,
2134            ..Default::default()
2135        });
2136
2137        right_levels.l0.sub_levels.push(Level {
2138            level_idx: 0,
2139            table_infos: vec![gen_sst_info(
2140                3,
2141                vec![5],
2142                FullKey::for_test(
2143                    TableId::new(5),
2144                    gen_key_from_str(VirtualNode::from_index(1), "1"),
2145                    0,
2146                )
2147                .encode()
2148                .into(),
2149                FullKey::for_test(
2150                    TableId::new(5),
2151                    gen_key_from_str(VirtualNode::from_index(200), "1"),
2152                    0,
2153                )
2154                .encode()
2155                .into(),
2156            )],
2157            sub_level_id: 103,
2158            level_type: LevelType::Nonoverlapping,
2159            total_file_size: 100,
2160            ..Default::default()
2161        });
2162
2163        {
2164            // test empty
2165            let mut left_levels = Levels::default();
2166            let right_levels = Levels::default();
2167
2168            group_split::merge_levels(&mut left_levels, right_levels);
2169        }
2170
2171        {
2172            // test empty left
2173            let mut left_levels = build_initial_compaction_group_levels(
2174                1,
2175                &CompactionConfig {
2176                    max_level: 6,
2177                    ..Default::default()
2178                },
2179            );
2180            let right_levels = right_levels.clone();
2181
2182            group_split::merge_levels(&mut left_levels, right_levels);
2183
2184            assert!(left_levels.l0.sub_levels.len() == 3);
2185            assert!(left_levels.l0.sub_levels[0].sub_level_id == 101);
2186            assert_eq!(100, left_levels.l0.sub_levels[0].total_file_size);
2187            assert!(left_levels.l0.sub_levels[1].sub_level_id == 102);
2188            assert_eq!(100, left_levels.l0.sub_levels[1].total_file_size);
2189            assert!(left_levels.l0.sub_levels[2].sub_level_id == 103);
2190            assert_eq!(100, left_levels.l0.sub_levels[2].total_file_size);
2191
2192            assert!(left_levels.levels[0].level_idx == 1);
2193            assert_eq!(300, left_levels.levels[0].total_file_size);
2194        }
2195
2196        {
2197            // test empty right
2198            let mut left_levels = left_levels.clone();
2199            let right_levels = build_initial_compaction_group_levels(
2200                2,
2201                &CompactionConfig {
2202                    max_level: 6,
2203                    ..Default::default()
2204                },
2205            );
2206
2207            group_split::merge_levels(&mut left_levels, right_levels);
2208
2209            assert!(left_levels.l0.sub_levels.len() == 3);
2210            assert!(left_levels.l0.sub_levels[0].sub_level_id == 101);
2211            assert_eq!(100, left_levels.l0.sub_levels[0].total_file_size);
2212            assert!(left_levels.l0.sub_levels[1].sub_level_id == 103);
2213            assert_eq!(100, left_levels.l0.sub_levels[1].total_file_size);
2214            assert!(left_levels.l0.sub_levels[2].sub_level_id == 105);
2215            assert_eq!(100, left_levels.l0.sub_levels[2].total_file_size);
2216
2217            assert!(left_levels.levels[0].level_idx == 1);
2218            assert_eq!(300, left_levels.levels[0].total_file_size);
2219        }
2220
2221        {
2222            let mut left_levels = left_levels.clone();
2223            let right_levels = right_levels.clone();
2224
2225            group_split::merge_levels(&mut left_levels, right_levels);
2226
2227            assert!(left_levels.l0.sub_levels.len() == 6);
2228            assert!(left_levels.l0.sub_levels[0].sub_level_id == 101);
2229            assert_eq!(100, left_levels.l0.sub_levels[0].total_file_size);
2230            assert!(left_levels.l0.sub_levels[1].sub_level_id == 103);
2231            assert_eq!(100, left_levels.l0.sub_levels[1].total_file_size);
2232            assert!(left_levels.l0.sub_levels[2].sub_level_id == 105);
2233            assert_eq!(100, left_levels.l0.sub_levels[2].total_file_size);
2234            assert!(left_levels.l0.sub_levels[3].sub_level_id == 106);
2235            assert_eq!(100, left_levels.l0.sub_levels[3].total_file_size);
2236            assert!(left_levels.l0.sub_levels[4].sub_level_id == 107);
2237            assert_eq!(100, left_levels.l0.sub_levels[4].total_file_size);
2238            assert!(left_levels.l0.sub_levels[5].sub_level_id == 108);
2239            assert_eq!(100, left_levels.l0.sub_levels[5].total_file_size);
2240
2241            assert!(left_levels.levels[0].level_idx == 1);
2242            assert_eq!(600, left_levels.levels[0].total_file_size);
2243        }
2244    }
2245
2246    #[test]
2247    fn test_get_split_pos() {
2248        let epoch = test_epoch(1);
2249        let s1 = gen_sstable_info(1, vec![1, 2], epoch);
2250        let s2 = gen_sstable_info(2, vec![3, 4, 5], epoch);
2251        let s3 = gen_sstable_info(3, vec![6, 7], epoch);
2252
2253        let ssts = vec![s1, s2, s3];
2254        let split_key = group_split::build_split_key(4.into(), VirtualNode::ZERO);
2255
2256        let pos = group_split::get_split_pos(&ssts, split_key.clone());
2257        assert_eq!(1, pos);
2258
2259        let pos = group_split::get_split_pos(&vec![], split_key);
2260        assert_eq!(0, pos);
2261    }
2262
2263    #[test]
2264    fn test_split_sst() {
2265        let epoch = test_epoch(1);
2266        let sst = gen_sstable_info(1, vec![1, 2, 3, 5], epoch);
2267
2268        {
2269            let split_key = group_split::build_split_key(3.into(), VirtualNode::ZERO);
2270            let origin_sst = sst.clone();
2271            let sst_size = origin_sst.sst_size;
2272            let split_type = group_split::need_to_split(&origin_sst, split_key.clone());
2273            assert_eq!(SstSplitType::Both, split_type);
2274
2275            let mut new_sst_id = 10.into();
2276            let (origin_sst, branched_sst) = group_split::split_sst(
2277                origin_sst,
2278                &mut new_sst_id,
2279                split_key,
2280                sst_size / 2,
2281                sst_size / 2,
2282            );
2283
2284            let origin_sst = origin_sst.unwrap();
2285            let branched_sst = branched_sst.unwrap();
2286
2287            assert!(origin_sst.key_range.right_exclusive);
2288            assert!(
2289                origin_sst
2290                    .key_range
2291                    .right
2292                    .cmp(&branched_sst.key_range.left)
2293                    .is_le()
2294            );
2295            assert!(origin_sst.table_ids.is_sorted());
2296            assert!(branched_sst.table_ids.is_sorted());
2297            assert!(origin_sst.table_ids.last().unwrap() < branched_sst.table_ids.first().unwrap());
2298            assert!(branched_sst.sst_size < origin_sst.file_size);
2299            assert_eq!(10, branched_sst.sst_id);
2300            assert_eq!(11, origin_sst.sst_id);
2301            assert_eq!(3, branched_sst.table_ids.first().unwrap().as_raw_id()); // split table_id to right
2302        }
2303
2304        {
2305            // test un-exist table_id
2306            let split_key = group_split::build_split_key(4.into(), VirtualNode::ZERO);
2307            let origin_sst = sst.clone();
2308            let sst_size = origin_sst.sst_size;
2309            let split_type = group_split::need_to_split(&origin_sst, split_key.clone());
2310            assert_eq!(SstSplitType::Both, split_type);
2311
2312            let mut new_sst_id = 10.into();
2313            let (origin_sst, branched_sst) = group_split::split_sst(
2314                origin_sst,
2315                &mut new_sst_id,
2316                split_key,
2317                sst_size / 2,
2318                sst_size / 2,
2319            );
2320
2321            let origin_sst = origin_sst.unwrap();
2322            let branched_sst = branched_sst.unwrap();
2323
2324            assert!(origin_sst.key_range.right_exclusive);
2325            assert!(origin_sst.key_range.right.le(&branched_sst.key_range.left));
2326            assert!(origin_sst.table_ids.is_sorted());
2327            assert!(branched_sst.table_ids.is_sorted());
2328            assert!(origin_sst.table_ids.last().unwrap() < branched_sst.table_ids.first().unwrap());
2329            assert!(branched_sst.sst_size < origin_sst.file_size);
2330            assert_eq!(10, branched_sst.sst_id);
2331            assert_eq!(11, origin_sst.sst_id);
2332            assert_eq!(5, branched_sst.table_ids.first().unwrap().as_raw_id()); // split table_id to right
2333        }
2334
2335        {
2336            let split_key = group_split::build_split_key(6.into(), VirtualNode::ZERO);
2337            let split_type = group_split::need_to_split(&sst, split_key);
2338            assert_eq!(SstSplitType::Left, split_type);
2339        }
2340
2341        {
2342            let split_key = group_split::build_split_key(4.into(), VirtualNode::ZERO);
2343            let origin_sst = sst.clone();
2344            let split_type = group_split::need_to_split(&origin_sst, split_key);
2345            assert_eq!(SstSplitType::Both, split_type);
2346
2347            let split_key = group_split::build_split_key(1.into(), VirtualNode::ZERO);
2348            let origin_sst = sst;
2349            let split_type = group_split::need_to_split(&origin_sst, split_key);
2350            assert_eq!(SstSplitType::Right, split_type);
2351        }
2352
2353        {
2354            // test key_range left = right
2355            let mut sst = gen_sstable_info_impl(1, vec![1], epoch);
2356            sst.key_range.right = sst.key_range.left.clone();
2357            let sst: SstableInfo = sst.into();
2358            let split_key = group_split::build_split_key(1.into(), VirtualNode::ZERO);
2359            let origin_sst = sst;
2360            let sst_size = origin_sst.sst_size;
2361
2362            let mut new_sst_id = 10.into();
2363            let (origin_sst, branched_sst) = group_split::split_sst(
2364                origin_sst,
2365                &mut new_sst_id,
2366                split_key,
2367                sst_size / 2,
2368                sst_size / 2,
2369            );
2370
2371            assert!(origin_sst.is_none());
2372            assert!(branched_sst.is_some());
2373        }
2374    }
2375
2376    #[test]
2377    fn test_split_sst_info_for_level() {
2378        let mut version = HummockVersion {
2379            id: HummockVersionId::new(0),
2380            levels: HashMap::from_iter([(
2381                1.into(),
2382                build_initial_compaction_group_levels(
2383                    1,
2384                    &CompactionConfig {
2385                        max_level: 6,
2386                        ..Default::default()
2387                    },
2388                ),
2389            )]),
2390            ..Default::default()
2391        };
2392
2393        let cg1 = version.levels.get_mut(&1).unwrap();
2394
2395        cg1.levels[0] = Level {
2396            level_idx: 1,
2397            level_type: LevelType::Nonoverlapping,
2398            table_infos: vec![
2399                gen_sst_info(
2400                    1,
2401                    vec![3],
2402                    FullKey::for_test(
2403                        TableId::new(3),
2404                        gen_key_from_str(VirtualNode::from_index(1), "1"),
2405                        0,
2406                    )
2407                    .encode()
2408                    .into(),
2409                    FullKey::for_test(
2410                        TableId::new(3),
2411                        gen_key_from_str(VirtualNode::from_index(200), "1"),
2412                        0,
2413                    )
2414                    .encode()
2415                    .into(),
2416                ),
2417                gen_sst_info(
2418                    10,
2419                    vec![3, 4],
2420                    FullKey::for_test(
2421                        TableId::new(3),
2422                        gen_key_from_str(VirtualNode::from_index(201), "1"),
2423                        0,
2424                    )
2425                    .encode()
2426                    .into(),
2427                    FullKey::for_test(
2428                        TableId::new(4),
2429                        gen_key_from_str(VirtualNode::from_index(10), "1"),
2430                        0,
2431                    )
2432                    .encode()
2433                    .into(),
2434                ),
2435                gen_sst_info(
2436                    11,
2437                    vec![4],
2438                    FullKey::for_test(
2439                        TableId::new(4),
2440                        gen_key_from_str(VirtualNode::from_index(11), "1"),
2441                        0,
2442                    )
2443                    .encode()
2444                    .into(),
2445                    FullKey::for_test(
2446                        TableId::new(4),
2447                        gen_key_from_str(VirtualNode::from_index(200), "1"),
2448                        0,
2449                    )
2450                    .encode()
2451                    .into(),
2452                ),
2453            ],
2454            total_file_size: 300,
2455            ..Default::default()
2456        };
2457
2458        cg1.l0.sub_levels.push(Level {
2459            level_idx: 0,
2460            table_infos: vec![
2461                gen_sst_info(
2462                    2,
2463                    vec![2],
2464                    FullKey::for_test(
2465                        TableId::new(0),
2466                        gen_key_from_str(VirtualNode::from_index(1), "1"),
2467                        0,
2468                    )
2469                    .encode()
2470                    .into(),
2471                    FullKey::for_test(
2472                        TableId::new(2),
2473                        gen_key_from_str(VirtualNode::from_index(200), "1"),
2474                        0,
2475                    )
2476                    .encode()
2477                    .into(),
2478                ),
2479                gen_sst_info(
2480                    22,
2481                    vec![2],
2482                    FullKey::for_test(
2483                        TableId::new(0),
2484                        gen_key_from_str(VirtualNode::from_index(1), "1"),
2485                        0,
2486                    )
2487                    .encode()
2488                    .into(),
2489                    FullKey::for_test(
2490                        TableId::new(2),
2491                        gen_key_from_str(VirtualNode::from_index(200), "1"),
2492                        0,
2493                    )
2494                    .encode()
2495                    .into(),
2496                ),
2497                gen_sst_info(
2498                    23,
2499                    vec![2],
2500                    FullKey::for_test(
2501                        TableId::new(0),
2502                        gen_key_from_str(VirtualNode::from_index(1), "1"),
2503                        0,
2504                    )
2505                    .encode()
2506                    .into(),
2507                    FullKey::for_test(
2508                        TableId::new(2),
2509                        gen_key_from_str(VirtualNode::from_index(200), "1"),
2510                        0,
2511                    )
2512                    .encode()
2513                    .into(),
2514                ),
2515                gen_sst_info(
2516                    24,
2517                    vec![2],
2518                    FullKey::for_test(
2519                        TableId::new(2),
2520                        gen_key_from_str(VirtualNode::from_index(1), "1"),
2521                        0,
2522                    )
2523                    .encode()
2524                    .into(),
2525                    FullKey::for_test(
2526                        TableId::new(2),
2527                        gen_key_from_str(VirtualNode::from_index(200), "1"),
2528                        0,
2529                    )
2530                    .encode()
2531                    .into(),
2532                ),
2533                gen_sst_info(
2534                    25,
2535                    vec![2],
2536                    FullKey::for_test(
2537                        TableId::new(0),
2538                        gen_key_from_str(VirtualNode::from_index(1), "1"),
2539                        0,
2540                    )
2541                    .encode()
2542                    .into(),
2543                    FullKey::for_test(
2544                        TableId::new(0),
2545                        gen_key_from_str(VirtualNode::from_index(200), "1"),
2546                        0,
2547                    )
2548                    .encode()
2549                    .into(),
2550                ),
2551            ],
2552            sub_level_id: 101,
2553            level_type: LevelType::Overlapping,
2554            total_file_size: 300,
2555            ..Default::default()
2556        });
2557
2558        {
2559            // split Overlapping level
2560            let split_key = group_split::build_split_key(1.into(), VirtualNode::ZERO);
2561
2562            let mut new_sst_id = 100.into();
2563            let x = group_split::split_sst_info_for_level_v2(
2564                &mut cg1.l0.sub_levels[0],
2565                &mut new_sst_id,
2566                split_key,
2567            );
2568            // assert_eq!(3, x.len());
2569            // assert_eq!(100, x[0].sst_id);
2570            // assert_eq!(100, x[0].sst_size);
2571            // assert_eq!(101, x[1].sst_id);
2572            // assert_eq!(100, x[1].sst_size);
2573            // assert_eq!(102, x[2].sst_id);
2574            // assert_eq!(100, x[2].sst_size);
2575
2576            let mut right_l0 = OverlappingLevel {
2577                sub_levels: vec![],
2578                total_file_size: 0,
2579                uncompressed_file_size: 0,
2580            };
2581
2582            right_l0.sub_levels.push(Level {
2583                level_idx: 0,
2584                table_infos: x,
2585                sub_level_id: 101,
2586                total_file_size: 100,
2587                level_type: LevelType::Overlapping,
2588                ..Default::default()
2589            });
2590
2591            let right_levels = Levels {
2592                levels: vec![],
2593                l0: right_l0,
2594                ..Default::default()
2595            };
2596
2597            merge_levels(cg1, right_levels);
2598        }
2599
2600        {
2601            // test split empty level
2602            let mut new_sst_id = 100.into();
2603            let split_key = group_split::build_split_key(1.into(), VirtualNode::ZERO);
2604            let x = group_split::split_sst_info_for_level_v2(
2605                &mut cg1.levels[2],
2606                &mut new_sst_id,
2607                split_key,
2608            );
2609
2610            assert!(x.is_empty());
2611        }
2612
2613        {
2614            // test split to right Nonoverlapping level
2615            let mut cg1 = cg1.clone();
2616            let split_key = group_split::build_split_key(1.into(), VirtualNode::ZERO);
2617
2618            let mut new_sst_id = 100.into();
2619            let x = group_split::split_sst_info_for_level_v2(
2620                &mut cg1.levels[0],
2621                &mut new_sst_id,
2622                split_key,
2623            );
2624
2625            assert_eq!(3, x.len());
2626            assert_eq!(1, x[0].sst_id);
2627            assert_eq!(100, x[0].sst_size);
2628            assert_eq!(10, x[1].sst_id);
2629            assert_eq!(100, x[1].sst_size);
2630            assert_eq!(11, x[2].sst_id);
2631            assert_eq!(100, x[2].sst_size);
2632
2633            assert_eq!(0, cg1.levels[0].table_infos.len());
2634        }
2635
2636        {
2637            // test split to left Nonoverlapping level
2638            let mut cg1 = cg1.clone();
2639            let split_key = group_split::build_split_key(5.into(), VirtualNode::ZERO);
2640
2641            let mut new_sst_id = 100.into();
2642            let x = group_split::split_sst_info_for_level_v2(
2643                &mut cg1.levels[0],
2644                &mut new_sst_id,
2645                split_key,
2646            );
2647
2648            assert_eq!(0, x.len());
2649            assert_eq!(3, cg1.levels[0].table_infos.len());
2650        }
2651
2652        // {
2653        //     // test split to both Nonoverlapping level
2654        //     let mut cg1 = cg1.clone();
2655        //     let split_key = build_split_key(3, VirtualNode::MAX);
2656
2657        //     let mut new_sst_id = 100;
2658        //     let x = group_split::split_sst_info_for_level_v2(
2659        //         &mut cg1.levels[0],
2660        //         &mut new_sst_id,
2661        //         split_key,
2662        //     );
2663
2664        //     assert_eq!(2, x.len());
2665        //     assert_eq!(100, x[0].sst_id);
2666        //     assert_eq!(100 / 2, x[0].sst_size);
2667        //     assert_eq!(11, x[1].sst_id);
2668        //     assert_eq!(100, x[1].sst_size);
2669        //     assert_eq!(vec![3, 4], x[0].table_ids);
2670
2671        //     assert_eq!(2, cg1.levels[0].table_infos.len());
2672        //     assert_eq!(101, cg1.levels[0].table_infos[1].sst_id);
2673        //     assert_eq!(100 / 2, cg1.levels[0].table_infos[1].sst_size);
2674        //     assert_eq!(vec![3], cg1.levels[0].table_infos[1].table_ids);
2675        // }
2676
2677        {
2678            // test split to both Nonoverlapping level
2679            let mut cg1 = cg1.clone();
2680            let split_key = group_split::build_split_key(4.into(), VirtualNode::ZERO);
2681
2682            let mut new_sst_id = 100.into();
2683            let x = group_split::split_sst_info_for_level_v2(
2684                &mut cg1.levels[0],
2685                &mut new_sst_id,
2686                split_key,
2687            );
2688
2689            assert_eq!(2, x.len());
2690            assert_eq!(100, x[0].sst_id);
2691            assert_eq!(100 / 2, x[0].sst_size);
2692            assert_eq!(11, x[1].sst_id);
2693            assert_eq!(100, x[1].sst_size);
2694            assert_eq!(vec![TableId::new(4)], x[1].table_ids);
2695
2696            assert_eq!(2, cg1.levels[0].table_infos.len());
2697            assert_eq!(101, cg1.levels[0].table_infos[1].sst_id);
2698            assert_eq!(100 / 2, cg1.levels[0].table_infos[1].sst_size);
2699            assert_eq!(
2700                vec![TableId::new(3)],
2701                cg1.levels[0].table_infos[1].table_ids
2702            );
2703        }
2704    }
2705
2706    #[test]
2707    fn test_validate_version_l0_sublevel_order() {
2708        let group_id = 1.into();
2709        let mut version = HummockVersion {
2710            levels: HashMap::from([(
2711                group_id,
2712                Levels {
2713                    group_id,
2714                    ..Default::default()
2715                },
2716            )]),
2717            ..Default::default()
2718        };
2719        for sub_level_id in [1, 2] {
2720            super::insert_new_sub_level(
2721                &mut version.levels.get_mut(&group_id).unwrap().l0,
2722                sub_level_id,
2723                LevelType::Overlapping,
2724                vec![gen_sstable_info(sub_level_id, vec![1], test_epoch(1))],
2725                None,
2726            );
2727        }
2728        let violations = super::validate_version(&version);
2729        assert!(violations.is_empty(), "{violations:?}");
2730
2731        version
2732            .levels
2733            .get_mut(&group_id)
2734            .unwrap()
2735            .l0
2736            .sub_levels
2737            .swap(0, 1);
2738        assert_eq!(
2739            super::validate_version(&version),
2740            vec!["GROUP 1 LEVEL 0: sub_level_id 1 <= prev_sub_level 2"],
2741        );
2742
2743        version.levels.get_mut(&group_id).unwrap().l0.sub_levels[0].sub_level_id = 1;
2744        assert_eq!(
2745            super::validate_version(&version),
2746            vec!["GROUP 1 LEVEL 0: sub_level_id 1 <= prev_sub_level 1"],
2747        );
2748    }
2749
2750    fn make_sst(sst_id: u64, table_ids: Vec<u32>, sst_size: u64) -> SstableInfo {
2751        SstableInfoInner {
2752            sst_id: sst_id.into(),
2753            object_id: sst_id.into(),
2754            table_ids: table_ids.into_iter().map(TableId::new).collect(),
2755            file_size: sst_size,
2756            sst_size,
2757            uncompressed_file_size: sst_size * 2,
2758            ..Default::default()
2759        }
2760        .into()
2761    }
2762
2763    #[test]
2764    fn test_level_normalize() {
2765        // Mixed: some SSTs empty, some not.
2766        let mut level = Level {
2767            level_idx: 1,
2768            level_type: LevelType::Nonoverlapping,
2769            table_infos: vec![
2770                make_sst(1, vec![1, 2], 100),
2771                make_sst(2, vec![], 200), // empty → removed
2772                make_sst(3, vec![3], 300),
2773                make_sst(4, vec![], 400), // empty → removed
2774            ],
2775            total_file_size: 9999, // intentionally wrong
2776            uncompressed_file_size: 9999,
2777            ..Default::default()
2778        };
2779
2780        level.normalize();
2781
2782        assert_eq!(level.table_infos.len(), 2);
2783        assert_eq!(1, level.table_infos[0].sst_id);
2784        assert_eq!(3, level.table_infos[1].sst_id);
2785        assert_eq!(level.total_file_size, 400);
2786        assert_eq!(level.uncompressed_file_size, 800);
2787
2788        // No empty SSTs: just recompute sizes.
2789        level.total_file_size = 0;
2790        level.uncompressed_file_size = 0;
2791        level.normalize();
2792        assert_eq!(level.table_infos.len(), 2);
2793        assert_eq!(level.total_file_size, 400);
2794        assert_eq!(level.uncompressed_file_size, 800);
2795
2796        // All empty → level becomes empty.
2797        level.table_infos = vec![make_sst(10, vec![], 100), make_sst(11, vec![], 200)];
2798        level.normalize();
2799        assert!(level.table_infos.is_empty());
2800        assert_eq!(level.total_file_size, 0);
2801        assert_eq!(level.uncompressed_file_size, 0);
2802    }
2803
2804    #[test]
2805    fn test_level_delete_ssts() {
2806        let mut level = Level {
2807            level_idx: 1,
2808            level_type: LevelType::Nonoverlapping,
2809            table_infos: vec![
2810                make_sst(1, vec![1], 100),
2811                make_sst(2, vec![2], 200),
2812                make_sst(3, vec![3], 300),
2813            ],
2814            total_file_size: 600,
2815            uncompressed_file_size: 1200,
2816            ..Default::default()
2817        };
2818
2819        let delete_ids: HashSet<crate::HummockSstableId> = HashSet::from([2.into()]);
2820        let changed = level.delete_ssts(&delete_ids);
2821
2822        assert!(changed);
2823        assert_eq!(level.table_infos.len(), 2);
2824        assert_eq!(1, level.table_infos[0].sst_id);
2825        assert_eq!(3, level.table_infos[1].sst_id);
2826        assert_eq!(level.total_file_size, 400);
2827        assert_eq!(level.uncompressed_file_size, 800);
2828
2829        // Delete non-existent id → no change, returns false.
2830        let delete_ids: HashSet<crate::HummockSstableId> = HashSet::from([999.into()]);
2831        let changed = level.delete_ssts(&delete_ids);
2832        assert!(!changed);
2833        assert_eq!(level.table_infos.len(), 2);
2834    }
2835
2836    #[test]
2837    fn test_level_prune_table_ids_from_ssts() {
2838        let mut level = Level {
2839            level_idx: 1,
2840            level_type: LevelType::Nonoverlapping,
2841            table_infos: vec![
2842                make_sst(1, vec![1, 2], 100),
2843                make_sst(2, vec![2, 3], 200),
2844                make_sst(3, vec![2], 300),
2845            ],
2846            total_file_size: 600,
2847            uncompressed_file_size: 1200,
2848            ..Default::default()
2849        };
2850
2851        let pruned_table_ids = HashSet::from([TableId::new(2)]);
2852        level.prune_table_ids_from_ssts(&pruned_table_ids);
2853
2854        // SST 3 should be removed (was only table_id=2)
2855        assert_eq!(level.table_infos.len(), 2);
2856        assert_eq!(level.table_infos[0].table_ids, vec![TableId::new(1)]);
2857        assert_eq!(level.table_infos[1].table_ids, vec![TableId::new(3)]);
2858        assert_eq!(level.total_file_size, 100 + 200);
2859        assert_eq!(level.uncompressed_file_size, 200 + 400);
2860
2861        // Prune remaining tables, so all SSTs are removed.
2862        let pruned_table_ids = HashSet::from([TableId::new(1), TableId::new(3)]);
2863        level.prune_table_ids_from_ssts(&pruned_table_ids);
2864        assert!(level.table_infos.is_empty());
2865        assert_eq!(level.total_file_size, 0);
2866        assert_eq!(level.uncompressed_file_size, 0);
2867    }
2868
2869    #[test]
2870    fn test_overlapping_level_normalize() {
2871        let mut l0 = OverlappingLevel {
2872            sub_levels: vec![
2873                Level {
2874                    level_idx: 0,
2875                    table_infos: vec![make_sst(1, vec![1], 100)],
2876                    total_file_size: 100,
2877                    uncompressed_file_size: 200,
2878                    sub_level_id: 1,
2879                    ..Default::default()
2880                },
2881                Level {
2882                    level_idx: 0,
2883                    table_infos: vec![], // empty → should be removed
2884                    total_file_size: 0,
2885                    uncompressed_file_size: 0,
2886                    sub_level_id: 2,
2887                    ..Default::default()
2888                },
2889                Level {
2890                    level_idx: 0,
2891                    table_infos: vec![make_sst(3, vec![3], 300)],
2892                    total_file_size: 300,
2893                    uncompressed_file_size: 600,
2894                    sub_level_id: 3,
2895                    ..Default::default()
2896                },
2897            ],
2898            total_file_size: 9999, // intentionally wrong
2899            uncompressed_file_size: 9999,
2900        };
2901
2902        l0.normalize();
2903
2904        assert_eq!(l0.sub_levels.len(), 2);
2905        assert_eq!(l0.sub_levels[0].sub_level_id, 1);
2906        assert_eq!(l0.sub_levels[1].sub_level_id, 3);
2907        assert_eq!(l0.total_file_size, 100 + 300);
2908        assert_eq!(l0.uncompressed_file_size, 200 + 600);
2909
2910        // All sub-levels empty → normalize clears everything.
2911        l0.sub_levels = vec![Level {
2912            level_idx: 0,
2913            table_infos: vec![],
2914            ..Default::default()
2915        }];
2916        l0.normalize();
2917        assert!(l0.sub_levels.is_empty());
2918        assert_eq!(l0.total_file_size, 0);
2919        assert_eq!(l0.uncompressed_file_size, 0);
2920    }
2921
2922    #[test]
2923    fn test_levels_prune_table_ids_from_ssts() {
2924        #[expect(deprecated)]
2925        let mut levels = Levels {
2926            l0: OverlappingLevel {
2927                sub_levels: vec![
2928                    Level {
2929                        level_idx: 0,
2930                        table_infos: vec![
2931                            make_sst(1, vec![10], 100), // table 10
2932                            make_sst(2, vec![20], 200), // table 20
2933                        ],
2934                        total_file_size: 300,
2935                        uncompressed_file_size: 600,
2936                        sub_level_id: 1,
2937                        ..Default::default()
2938                    },
2939                    Level {
2940                        level_idx: 0,
2941                        table_infos: vec![
2942                            make_sst(3, vec![10], 150), // table 10
2943                        ],
2944                        total_file_size: 150,
2945                        uncompressed_file_size: 300,
2946                        sub_level_id: 2,
2947                        ..Default::default()
2948                    },
2949                ],
2950                total_file_size: 450,
2951                uncompressed_file_size: 900,
2952            },
2953            levels: vec![Level {
2954                level_idx: 1,
2955                level_type: LevelType::Nonoverlapping,
2956                table_infos: vec![
2957                    make_sst(4, vec![10, 20], 400), // shared SST
2958                    make_sst(5, vec![10], 500),     // table 10 only
2959                ],
2960                total_file_size: 900,
2961                uncompressed_file_size: 1800,
2962                ..Default::default()
2963            }],
2964            group_id: 1.into(),
2965            parent_group_id: 0.into(),
2966            member_table_ids: vec![],
2967            compaction_group_version_id: 0,
2968        };
2969
2970        // Prune table 10 from SST metadata.
2971        levels.prune_table_ids_from_ssts(&HashSet::from([TableId::new(10)]));
2972
2973        assert_eq!(levels.l0.sub_levels.len(), 1);
2974        assert_eq!(levels.l0.sub_levels[0].sub_level_id, 1);
2975        assert_eq!(levels.l0.sub_levels[0].table_infos.len(), 1);
2976        assert_eq!(2, levels.l0.sub_levels[0].table_infos[0].sst_id);
2977        assert_eq!(levels.l0.sub_levels[0].total_file_size, 200);
2978        assert_eq!(levels.l0.sub_levels[0].uncompressed_file_size, 400);
2979
2980        assert_eq!(levels.l0.total_file_size, 200);
2981        assert_eq!(levels.l0.uncompressed_file_size, 400);
2982
2983        assert_eq!(levels.levels[0].table_infos.len(), 1);
2984        assert_eq!(4, levels.levels[0].table_infos[0].sst_id);
2985        assert_eq!(
2986            levels.levels[0].table_infos[0].table_ids,
2987            vec![TableId::new(20)]
2988        );
2989        assert_eq!(levels.levels[0].total_file_size, 400);
2990        assert_eq!(levels.levels[0].uncompressed_file_size, 800);
2991
2992        assert_eq!(levels.compaction_group_version_id, 1);
2993    }
2994
2995    #[test]
2996    fn test_apply_version_delta_prune_table_ids_from_ssts() {
2997        let mut version = HummockVersion {
2998            id: HummockVersionId::new(0),
2999            levels: HashMap::from_iter([(1.into(), {
3000                #[expect(deprecated)]
3001                let levels = Levels {
3002                    l0: OverlappingLevel {
3003                        sub_levels: vec![
3004                            Level {
3005                                level_idx: 0,
3006                                level_type: LevelType::Overlapping,
3007                                table_infos: vec![
3008                                    make_sst(1, vec![100], 50), // only table 100
3009                                    make_sst(2, vec![200], 60), // only table 200
3010                                ],
3011                                total_file_size: 110,
3012                                uncompressed_file_size: 220,
3013                                sub_level_id: 1,
3014                                ..Default::default()
3015                            },
3016                            Level {
3017                                level_idx: 0,
3018                                level_type: LevelType::Overlapping,
3019                                table_infos: vec![
3020                                    make_sst(3, vec![100], 70), // only table 100
3021                                ],
3022                                total_file_size: 70,
3023                                uncompressed_file_size: 140,
3024                                sub_level_id: 2,
3025                                ..Default::default()
3026                            },
3027                        ],
3028                        total_file_size: 180,
3029                        uncompressed_file_size: 360,
3030                    },
3031                    levels: vec![Level {
3032                        level_idx: 1,
3033                        level_type: LevelType::Nonoverlapping,
3034                        table_infos: vec![
3035                            make_sst(4, vec![100, 200], 80), // shared
3036                            make_sst(5, vec![100], 90),      // only table 100
3037                        ],
3038                        total_file_size: 170,
3039                        uncompressed_file_size: 340,
3040                        ..Default::default()
3041                    }],
3042                    group_id: 1.into(),
3043                    parent_group_id: 0.into(),
3044                    member_table_ids: vec![],
3045                    compaction_group_version_id: 0,
3046                };
3047                levels
3048            })]),
3049            ..Default::default()
3050        };
3051
3052        let version_delta = HummockVersionDelta {
3053            id: HummockVersionId::new(1),
3054            group_deltas: HashMap::from_iter([(
3055                1.into(),
3056                GroupDeltas {
3057                    group_deltas: vec![GroupDelta::PruneTableIdsFromSsts(HashSet::from([
3058                        TableId::new(100),
3059                    ]))],
3060                },
3061            )]),
3062            ..Default::default()
3063        };
3064
3065        version.apply_version_delta(&version_delta);
3066
3067        let cg = version.get_compaction_group_levels(1.into());
3068
3069        assert_eq!(
3070            cg.l0.sub_levels.len(),
3071            1,
3072            "empty sub-level should be removed"
3073        );
3074        assert_eq!(cg.l0.sub_levels[0].sub_level_id, 1);
3075        assert_eq!(cg.l0.sub_levels[0].table_infos.len(), 1);
3076        assert_eq!(2, cg.l0.sub_levels[0].table_infos[0].sst_id);
3077        assert_eq!(
3078            cg.l0.sub_levels[0].table_infos[0].table_ids,
3079            vec![TableId::new(200)]
3080        );
3081
3082        assert_eq!(cg.l0.total_file_size, 60);
3083        assert_eq!(cg.l0.uncompressed_file_size, 120);
3084
3085        assert_eq!(cg.levels[0].table_infos.len(), 1);
3086        assert_eq!(4, cg.levels[0].table_infos[0].sst_id);
3087        assert_eq!(
3088            cg.levels[0].table_infos[0].table_ids,
3089            vec![TableId::new(200)]
3090        );
3091        assert_eq!(cg.levels[0].total_file_size, 80);
3092        assert_eq!(cg.levels[0].uncompressed_file_size, 160);
3093
3094        assert_eq!(cg.compaction_group_version_id, 1);
3095    }
3096
3097    #[test]
3098    fn test_apply_version_delta_removes_dropped_table_watermark() {
3099        let table_id = TableId::new(100);
3100        let mut version = HummockVersion {
3101            id: HummockVersionId::new(0),
3102            table_watermarks: HashMap::from([(
3103                table_id,
3104                Arc::new(TableWatermarks::single_epoch(
3105                    test_epoch(1),
3106                    vec![VnodeWatermark::new(
3107                        Arc::new(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)),
3108                        Bytes::from_static(b"watermark"),
3109                    )],
3110                    WatermarkDirection::Ascending,
3111                    WatermarkSerdeType::PkPrefix,
3112                )),
3113            )]),
3114            ..Default::default()
3115        };
3116
3117        let version_delta = HummockVersionDelta {
3118            id: HummockVersionId::new(1),
3119            prev_id: HummockVersionId::new(0),
3120            removed_table_ids: HashSet::from([table_id]),
3121            ..Default::default()
3122        };
3123
3124        version.apply_version_delta(&version_delta);
3125
3126        assert!(!version.table_watermarks.contains_key(&table_id));
3127    }
3128
3129    #[test]
3130    fn test_prune_stale_table_ids_from_ssts() {
3131        let live_table_id = TableId::new(100);
3132        let stale_table_id = TableId::new(200);
3133        let mut version = HummockVersion {
3134            id: HummockVersionId::new(0),
3135            levels: HashMap::from_iter([(1.into(), {
3136                #[expect(deprecated)]
3137                let levels = Levels {
3138                    l0: OverlappingLevel {
3139                        sub_levels: vec![Level {
3140                            level_idx: 0,
3141                            level_type: LevelType::Overlapping,
3142                            table_infos: vec![make_sst(
3143                                1,
3144                                vec![live_table_id.as_raw_id(), stale_table_id.as_raw_id()],
3145                                50,
3146                            )],
3147                            total_file_size: 50,
3148                            uncompressed_file_size: 100,
3149                            sub_level_id: 1,
3150                            ..Default::default()
3151                        }],
3152                        total_file_size: 50,
3153                        uncompressed_file_size: 100,
3154                    },
3155                    levels: vec![Level {
3156                        level_idx: 1,
3157                        level_type: LevelType::Nonoverlapping,
3158                        table_infos: vec![make_sst(2, vec![stale_table_id.as_raw_id()], 60)],
3159                        total_file_size: 60,
3160                        uncompressed_file_size: 120,
3161                        ..Default::default()
3162                    }],
3163                    group_id: 1.into(),
3164                    parent_group_id: 0.into(),
3165                    member_table_ids: vec![],
3166                    compaction_group_version_id: 0,
3167                };
3168                levels
3169            })]),
3170            state_table_info: HummockVersionStateTableInfo::from_protobuf_owned(
3171                HashMap::from_iter([(
3172                    live_table_id,
3173                    StateTableInfo {
3174                        committed_epoch: 1,
3175                        compaction_group_id: 1.into(),
3176                    },
3177                )]),
3178            ),
3179            ..Default::default()
3180        };
3181
3182        assert_eq!(version.prune_stale_table_ids_from_ssts(), 1);
3183
3184        let cg = version.get_compaction_group_levels(1.into());
3185        assert_eq!(cg.l0.sub_levels.len(), 1);
3186        assert_eq!(cg.l0.sub_levels[0].table_infos.len(), 1);
3187        assert_eq!(cg.l0.sub_levels[0].table_infos[0].sst_id, 1);
3188        assert_eq!(
3189            cg.l0.sub_levels[0].table_infos[0].table_ids,
3190            vec![live_table_id]
3191        );
3192        assert!(cg.levels[0].table_infos.is_empty());
3193        assert_eq!(cg.compaction_group_version_id, 1);
3194    }
3195
3196    #[test]
3197    fn test_prune_stale_table_ids_from_ssts_skips_legacy_member_table_ids() {
3198        let mut version = HummockVersion {
3199            id: HummockVersionId::new(0),
3200            levels: HashMap::from_iter([(1.into(), {
3201                #[expect(deprecated)]
3202                let levels = Levels {
3203                    l0: OverlappingLevel {
3204                        sub_levels: vec![Level {
3205                            level_idx: 0,
3206                            level_type: LevelType::Overlapping,
3207                            table_infos: vec![make_sst(1, vec![100, 200], 50)],
3208                            total_file_size: 50,
3209                            uncompressed_file_size: 100,
3210                            sub_level_id: 1,
3211                            ..Default::default()
3212                        }],
3213                        total_file_size: 50,
3214                        uncompressed_file_size: 100,
3215                    },
3216                    group_id: 1.into(),
3217                    parent_group_id: 0.into(),
3218                    member_table_ids: vec![100, 200],
3219                    compaction_group_version_id: 0,
3220                    ..Default::default()
3221                };
3222                levels
3223            })]),
3224            ..Default::default()
3225        };
3226
3227        assert_eq!(version.prune_stale_table_ids_from_ssts(), 0);
3228
3229        let cg = version.get_compaction_group_levels(1.into());
3230        assert_eq!(
3231            cg.l0.sub_levels[0].table_infos[0].table_ids,
3232            vec![TableId::new(100), TableId::new(200)]
3233        );
3234        assert_eq!(cg.compaction_group_version_id, 0);
3235    }
3236
3237    #[test]
3238    fn test_apply_version_delta_compact_l0() {
3239        let mut version = HummockVersion {
3240            id: HummockVersionId::new(0),
3241            levels: HashMap::from_iter([(1.into(), {
3242                #[expect(deprecated)]
3243                let levels = Levels {
3244                    l0: OverlappingLevel {
3245                        sub_levels: vec![
3246                            Level {
3247                                level_idx: 0,
3248                                level_type: LevelType::Nonoverlapping,
3249                                table_infos: vec![
3250                                    make_sst(1, vec![1], 100),
3251                                    make_sst(2, vec![2], 200),
3252                                ],
3253                                total_file_size: 300,
3254                                uncompressed_file_size: 600,
3255                                sub_level_id: 1,
3256                                ..Default::default()
3257                            },
3258                            Level {
3259                                level_idx: 0,
3260                                level_type: LevelType::Nonoverlapping,
3261                                table_infos: vec![make_sst(3, vec![3], 300)],
3262                                total_file_size: 300,
3263                                uncompressed_file_size: 600,
3264                                sub_level_id: 2,
3265                                ..Default::default()
3266                            },
3267                        ],
3268                        total_file_size: 600,
3269                        uncompressed_file_size: 1200,
3270                    },
3271                    levels: vec![Level {
3272                        level_idx: 1,
3273                        level_type: LevelType::Nonoverlapping,
3274                        table_infos: vec![],
3275                        total_file_size: 0,
3276                        uncompressed_file_size: 0,
3277                        ..Default::default()
3278                    }],
3279                    group_id: 1.into(),
3280                    parent_group_id: 0.into(),
3281                    member_table_ids: vec![],
3282                    compaction_group_version_id: 0,
3283                };
3284                levels
3285            })]),
3286            ..Default::default()
3287        };
3288
3289        let version_delta = HummockVersionDelta {
3290            id: HummockVersionId::new(1),
3291            group_deltas: HashMap::from_iter([(
3292                1.into(),
3293                GroupDeltas {
3294                    group_deltas: vec![
3295                        GroupDelta::IntraLevel(IntraLevelDelta::new(
3296                            0, // L0
3297                            0,
3298                            HashSet::from([1.into(), 2.into(), 3.into()]),
3299                            vec![],
3300                            0,
3301                            0,
3302                        )),
3303                        GroupDelta::IntraLevel(IntraLevelDelta::new(
3304                            1, // L1
3305                            0,
3306                            HashSet::new(),
3307                            vec![make_sst(10, vec![1, 2, 3], 500)],
3308                            0,
3309                            0,
3310                        )),
3311                    ],
3312                },
3313            )]),
3314            ..Default::default()
3315        };
3316
3317        version.apply_version_delta(&version_delta);
3318
3319        let cg = version.get_compaction_group_levels(1.into());
3320
3321        assert!(cg.l0.sub_levels.is_empty());
3322        assert_eq!(cg.l0.total_file_size, 0);
3323        assert_eq!(cg.l0.uncompressed_file_size, 0);
3324
3325        assert_eq!(cg.levels[0].table_infos.len(), 1);
3326        assert_eq!(10, cg.levels[0].table_infos[0].sst_id);
3327        assert_eq!(cg.levels[0].total_file_size, 500);
3328    }
3329}