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