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