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