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