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