1use std::collections::hash_map::Entry;
16use std::collections::{BTreeSet, HashMap, HashSet};
17use std::mem::{replace, size_of};
18use std::ops::Deref;
19use std::sync::{Arc, LazyLock};
20
21use itertools::Itertools;
22use risingwave_common::catalog::TableId;
23use risingwave_common::util::epoch::INVALID_EPOCH;
24use risingwave_pb::hummock::group_delta::{DeltaType, PbDeltaType};
25use risingwave_pb::hummock::hummock_version_delta::PbGroupDeltas;
26use risingwave_pb::hummock::*;
27use tracing::warn;
28
29use crate::compaction_group::StaticCompactionGroupId;
30use crate::compaction_group::hummock_version_ext::build_initial_compaction_group_levels;
31use crate::level::LevelsCommon;
32use crate::sstable_info::SstableInfo;
33use crate::table_watermark::TableWatermarks;
34use crate::vector_index::{VectorIndex, VectorIndexDelta};
35use crate::{
36 CompactionGroupId, FIRST_VERSION_ID, HummockEpoch, HummockObjectId, HummockSstableId,
37 HummockSstableObjectId, HummockVersionId,
38};
39
40pub const MAX_HUMMOCK_VERSION_ID: HummockVersionId = HummockVersionId::new(i64::MAX as _);
41
42#[derive(Debug, Clone, PartialEq)]
43pub struct HummockVersionStateTableInfo {
44 state_table_info: HashMap<TableId, PbStateTableInfo>,
45
46 compaction_group_member_tables: HashMap<CompactionGroupId, BTreeSet<TableId>>,
48}
49
50impl HummockVersionStateTableInfo {
51 pub fn empty() -> Self {
52 Self {
53 state_table_info: HashMap::new(),
54 compaction_group_member_tables: HashMap::new(),
55 }
56 }
57
58 fn build_compaction_group_member_tables(
59 state_table_info: &HashMap<TableId, PbStateTableInfo>,
60 ) -> HashMap<CompactionGroupId, BTreeSet<TableId>> {
61 let mut ret: HashMap<_, BTreeSet<_>> = HashMap::new();
62 for (table_id, info) in state_table_info {
63 assert!(
64 ret.entry(info.compaction_group_id)
65 .or_default()
66 .insert(*table_id)
67 );
68 }
69 ret
70 }
71
72 pub fn build_table_compaction_group_id(&self) -> HashMap<TableId, CompactionGroupId> {
73 self.state_table_info
74 .iter()
75 .map(|(table_id, info)| (*table_id, info.compaction_group_id))
76 .collect()
77 }
78
79 pub fn from_protobuf(state_table_info: &HashMap<TableId, PbStateTableInfo>) -> Self {
80 let state_table_info = state_table_info
81 .iter()
82 .map(|(table_id, info)| (*table_id, *info))
83 .collect();
84 let compaction_group_member_tables =
85 Self::build_compaction_group_member_tables(&state_table_info);
86 Self {
87 state_table_info,
88 compaction_group_member_tables,
89 }
90 }
91
92 pub fn from_protobuf_owned(state_table_info: HashMap<TableId, PbStateTableInfo>) -> Self {
93 let compaction_group_member_tables =
94 Self::build_compaction_group_member_tables(&state_table_info);
95 Self {
96 state_table_info,
97 compaction_group_member_tables,
98 }
99 }
100
101 pub fn apply_delta(
102 &mut self,
103 delta: &HashMap<TableId, StateTableInfoDelta>,
104 removed_table_id: &HashSet<TableId>,
105 ) -> (HashMap<TableId, Option<StateTableInfo>>, bool) {
106 let mut changed_table = HashMap::new();
107 let mut has_bumped_committed_epoch = false;
108 fn remove_table_from_compaction_group(
109 compaction_group_member_tables: &mut HashMap<CompactionGroupId, BTreeSet<TableId>>,
110 compaction_group_id: CompactionGroupId,
111 table_id: TableId,
112 ) {
113 let member_tables = compaction_group_member_tables
114 .get_mut(&compaction_group_id)
115 .expect("should exist");
116 assert!(member_tables.remove(&table_id));
117 if member_tables.is_empty() {
118 assert!(
119 compaction_group_member_tables
120 .remove(&compaction_group_id)
121 .is_some()
122 );
123 }
124 }
125 for table_id in removed_table_id {
126 if let Some(prev_info) = self.state_table_info.remove(table_id) {
127 remove_table_from_compaction_group(
128 &mut self.compaction_group_member_tables,
129 prev_info.compaction_group_id,
130 *table_id,
131 );
132 assert!(changed_table.insert(*table_id, Some(prev_info)).is_none());
133 } else {
134 warn!(
135 %table_id,
136 "table to remove does not exist"
137 );
138 }
139 }
140 for (table_id, delta) in delta {
141 if removed_table_id.contains(table_id) {
142 continue;
143 }
144 let new_info = StateTableInfo {
145 committed_epoch: delta.committed_epoch,
146 compaction_group_id: delta.compaction_group_id,
147 };
148 match self.state_table_info.entry(*table_id) {
149 Entry::Occupied(mut entry) => {
150 let prev_info = entry.get_mut();
151 assert!(
152 new_info.committed_epoch >= prev_info.committed_epoch,
153 "state table info regress. table id: {}, prev_info: {:?}, new_info: {:?}",
154 table_id,
155 prev_info,
156 new_info
157 );
158 if new_info.committed_epoch > prev_info.committed_epoch {
159 has_bumped_committed_epoch = true;
160 }
161 if prev_info.compaction_group_id != new_info.compaction_group_id {
162 remove_table_from_compaction_group(
164 &mut self.compaction_group_member_tables,
165 prev_info.compaction_group_id,
166 *table_id,
167 );
168 assert!(
169 self.compaction_group_member_tables
170 .entry(new_info.compaction_group_id)
171 .or_default()
172 .insert(*table_id)
173 );
174 }
175 let prev_info = replace(prev_info, new_info);
176 changed_table.insert(*table_id, Some(prev_info));
177 }
178 Entry::Vacant(entry) => {
179 assert!(
180 self.compaction_group_member_tables
181 .entry(new_info.compaction_group_id)
182 .or_default()
183 .insert(*table_id)
184 );
185 has_bumped_committed_epoch = true;
186 entry.insert(new_info);
187 changed_table.insert(*table_id, None);
188 }
189 }
190 }
191 debug_assert_eq!(
192 self.compaction_group_member_tables,
193 Self::build_compaction_group_member_tables(&self.state_table_info)
194 );
195 (changed_table, has_bumped_committed_epoch)
196 }
197
198 pub fn info(&self) -> &HashMap<TableId, StateTableInfo> {
199 &self.state_table_info
200 }
201
202 pub fn compaction_group_member_table_ids(
203 &self,
204 compaction_group_id: CompactionGroupId,
205 ) -> &BTreeSet<TableId> {
206 static EMPTY_SET: LazyLock<BTreeSet<TableId>> = LazyLock::new(BTreeSet::new);
207 self.compaction_group_member_tables
208 .get(&compaction_group_id)
209 .unwrap_or_else(|| EMPTY_SET.deref())
210 }
211
212 pub fn compaction_group_member_tables(&self) -> &HashMap<CompactionGroupId, BTreeSet<TableId>> {
213 &self.compaction_group_member_tables
214 }
215
216 pub fn max_table_committed_epoch(&self) -> Option<HummockEpoch> {
217 self.state_table_info
218 .values()
219 .map(|info| info.committed_epoch)
220 .max()
221 }
222}
223
224#[derive(Debug, Clone, PartialEq)]
225pub struct HummockVersionCommon<T> {
226 pub id: HummockVersionId,
227 pub levels: HashMap<CompactionGroupId, LevelsCommon<T>>,
228 #[deprecated]
229 pub(crate) max_committed_epoch: u64,
230 pub table_watermarks: HashMap<TableId, Arc<TableWatermarks>>,
231 pub state_table_info: HummockVersionStateTableInfo,
232 pub vector_indexes: HashMap<TableId, VectorIndex>,
233}
234
235pub type HummockVersion = HummockVersionCommon<SstableInfo>;
236
237impl Default for HummockVersion {
238 fn default() -> Self {
239 HummockVersion::from(&PbHummockVersion::default())
240 }
241}
242
243impl<T> HummockVersionCommon<T>
244where
245 T: for<'a> From<&'a PbSstableInfo>,
246 PbSstableInfo: for<'a> From<&'a T>,
247{
248 pub fn from_rpc_protobuf(pb_version: &PbHummockVersion) -> Self {
251 pb_version.into()
252 }
253
254 pub fn from_persisted_protobuf(pb_version: &PbHummockVersion) -> Self {
257 pb_version.into()
258 }
259
260 pub fn to_protobuf(&self) -> PbHummockVersion {
261 self.into()
262 }
263}
264
265impl<T> HummockVersionCommon<T>
266where
267 T: From<PbSstableInfo>,
268 PbSstableInfo: for<'a> From<&'a T>,
269{
270 pub fn from_persisted_protobuf_owned(pb_version: PbHummockVersion) -> Self {
273 pb_version.into()
274 }
275}
276
277impl HummockVersion {
278 pub fn estimated_encode_len(&self) -> usize {
279 self.levels.len() * size_of::<CompactionGroupId>()
280 + self
281 .levels
282 .values()
283 .map(|level| level.estimated_encode_len())
284 .sum::<usize>()
285 + self.table_watermarks.len() * size_of::<u32>()
286 + self
287 .table_watermarks
288 .values()
289 .map(|table_watermark| table_watermark.estimated_encode_len())
290 .sum::<usize>()
291 }
292}
293
294impl<T> From<&PbHummockVersion> for HummockVersionCommon<T>
295where
296 T: for<'a> From<&'a PbSstableInfo>,
297{
298 fn from(pb_version: &PbHummockVersion) -> Self {
299 #[expect(deprecated)]
300 Self {
301 id: pb_version.id,
302 levels: pb_version
303 .levels
304 .iter()
305 .map(|(group_id, levels)| (*group_id, LevelsCommon::from(levels)))
306 .collect(),
307 max_committed_epoch: pb_version.max_committed_epoch,
308 table_watermarks: pb_version
309 .table_watermarks
310 .iter()
311 .map(|(table_id, table_watermark)| {
312 (*table_id, Arc::new(TableWatermarks::from(table_watermark)))
313 })
314 .collect(),
315 state_table_info: HummockVersionStateTableInfo::from_protobuf(
316 &pb_version.state_table_info,
317 ),
318 vector_indexes: pb_version
319 .vector_indexes
320 .iter()
321 .map(|(table_id, index)| (*table_id, index.clone().into()))
322 .collect(),
323 }
324 }
325}
326
327impl<T> From<PbHummockVersion> for HummockVersionCommon<T>
328where
329 T: From<PbSstableInfo>,
330{
331 fn from(pb_version: PbHummockVersion) -> Self {
332 #[expect(deprecated)]
333 Self {
334 id: pb_version.id,
335 levels: pb_version
336 .levels
337 .into_iter()
338 .map(|(group_id, levels)| (group_id, LevelsCommon::from(levels)))
339 .collect(),
340 max_committed_epoch: pb_version.max_committed_epoch,
341 table_watermarks: pb_version
342 .table_watermarks
343 .into_iter()
344 .map(|(table_id, table_watermark)| {
345 (table_id, Arc::new(TableWatermarks::from(table_watermark)))
346 })
347 .collect(),
348 state_table_info: HummockVersionStateTableInfo::from_protobuf_owned(
349 pb_version.state_table_info,
350 ),
351 vector_indexes: pb_version
352 .vector_indexes
353 .into_iter()
354 .map(|(table_id, index)| (table_id, index.into()))
355 .collect(),
356 }
357 }
358}
359
360impl<T> From<&HummockVersionCommon<T>> for PbHummockVersion
361where
362 PbSstableInfo: for<'a> From<&'a T>,
363{
364 fn from(version: &HummockVersionCommon<T>) -> Self {
365 #[expect(deprecated)]
366 Self {
367 id: version.id,
368 levels: version
369 .levels
370 .iter()
371 .map(|(group_id, levels)| (*group_id, levels.into()))
372 .collect(),
373 max_committed_epoch: version.max_committed_epoch,
374 table_watermarks: version
375 .table_watermarks
376 .iter()
377 .map(|(table_id, watermark)| (*table_id, watermark.as_ref().into()))
378 .collect(),
379 table_change_logs: Default::default(),
380 state_table_info: version.state_table_info.state_table_info.clone(),
381 vector_indexes: version
382 .vector_indexes
383 .iter()
384 .map(|(table_id, index)| (*table_id, index.clone().into()))
385 .collect(),
386 }
387 }
388}
389
390impl<T> From<HummockVersionCommon<T>> for PbHummockVersion
391where
392 PbSstableInfo: From<T>,
393 PbSstableInfo: for<'a> From<&'a T>,
394{
395 fn from(version: HummockVersionCommon<T>) -> Self {
396 #[expect(deprecated)]
397 Self {
398 id: version.id,
399 levels: version
400 .levels
401 .into_iter()
402 .map(|(group_id, levels)| (group_id, levels.into()))
403 .collect(),
404 max_committed_epoch: version.max_committed_epoch,
405 table_watermarks: version
406 .table_watermarks
407 .into_iter()
408 .map(|(table_id, watermark)| (table_id, watermark.as_ref().into()))
409 .collect(),
410 table_change_logs: Default::default(),
411 state_table_info: version.state_table_info.state_table_info.clone(),
412 vector_indexes: version
413 .vector_indexes
414 .into_iter()
415 .map(|(table_id, index)| (table_id, index.into()))
416 .collect(),
417 }
418 }
419}
420
421impl HummockVersion {
422 pub fn next_version_id(&self) -> HummockVersionId {
423 self.id + 1
424 }
425
426 pub fn need_fill_backward_compatible_state_table_info_delta(&self) -> bool {
427 self.state_table_info.state_table_info.is_empty()
429 && self.levels.values().any(|group| {
430 #[expect(deprecated)]
432 !group.member_table_ids.is_empty()
433 })
434 }
435
436 pub fn may_fill_backward_compatible_state_table_info_delta(
437 &self,
438 delta: &mut HummockVersionDelta,
439 ) {
440 #[expect(deprecated)]
441 for (cg_id, group) in &self.levels {
443 for table_id in &group.member_table_ids {
444 assert!(
445 delta
446 .state_table_info_delta
447 .insert(
448 (*table_id).into(),
449 StateTableInfoDelta {
450 committed_epoch: self.max_committed_epoch,
451 compaction_group_id: *cg_id,
452 }
453 )
454 .is_none(),
455 "duplicate table id {} in cg {}",
456 table_id,
457 cg_id
458 );
459 }
460 }
461 }
462
463 pub fn create_init_version(default_compaction_config: Arc<CompactionConfig>) -> HummockVersion {
464 #[expect(deprecated)]
465 let mut init_version = HummockVersion {
466 id: FIRST_VERSION_ID,
467 levels: Default::default(),
468 max_committed_epoch: INVALID_EPOCH,
469 table_watermarks: HashMap::new(),
470 state_table_info: HummockVersionStateTableInfo::empty(),
471 vector_indexes: Default::default(),
472 };
473 for group_id in [
474 StaticCompactionGroupId::StateDefault as CompactionGroupId,
475 StaticCompactionGroupId::MaterializedView as CompactionGroupId,
476 ] {
477 init_version.levels.insert(
478 group_id,
479 build_initial_compaction_group_levels(group_id, default_compaction_config.as_ref()),
480 );
481 }
482 init_version
483 }
484
485 pub fn version_delta_after(&self) -> HummockVersionDelta {
486 #[expect(deprecated)]
487 HummockVersionDelta {
488 id: self.next_version_id(),
489 prev_id: self.id,
490 trivial_move: false,
491 max_committed_epoch: self.max_committed_epoch,
492 group_deltas: Default::default(),
493 new_table_watermarks: HashMap::new(),
494 removed_table_ids: HashSet::new(),
495 state_table_info_delta: Default::default(),
496 vector_index_delta: Default::default(),
497 }
498 }
499}
500
501impl<T> HummockVersionCommon<T> {
502 pub fn table_committed_epoch(&self, table_id: TableId) -> Option<u64> {
503 self.state_table_info
504 .info()
505 .get(&table_id)
506 .map(|info| info.committed_epoch)
507 }
508}
509
510#[derive(Debug, PartialEq, Clone)]
511pub struct HummockVersionDeltaCommon<T> {
512 pub id: HummockVersionId,
513 pub prev_id: HummockVersionId,
514 pub group_deltas: HashMap<CompactionGroupId, GroupDeltasCommon<T>>,
515 #[deprecated]
516 pub(crate) max_committed_epoch: u64,
517 pub trivial_move: bool,
518 pub new_table_watermarks: HashMap<TableId, TableWatermarks>,
519 pub removed_table_ids: HashSet<TableId>,
520 pub state_table_info_delta: HashMap<TableId, StateTableInfoDelta>,
521 pub vector_index_delta: HashMap<TableId, VectorIndexDelta>,
522}
523
524pub type HummockVersionDelta = HummockVersionDeltaCommon<SstableInfo>;
525
526impl Default for HummockVersionDelta {
527 fn default() -> Self {
528 HummockVersionDelta::from(&PbHummockVersionDelta::default())
529 }
530}
531
532impl<T> HummockVersionDeltaCommon<T>
533where
534 T: for<'a> From<&'a PbSstableInfo>,
535 PbSstableInfo: for<'a> From<&'a T>,
536{
537 pub fn from_persisted_protobuf(delta: &PbHummockVersionDelta) -> Self {
540 warn_if_legacy_change_log_delta_is_present(delta);
541 delta.into()
542 }
543
544 pub fn from_rpc_protobuf(delta: &PbHummockVersionDelta) -> Self {
547 delta.into()
548 }
549
550 pub fn to_protobuf(&self) -> PbHummockVersionDelta {
551 self.into()
552 }
553}
554
555impl<T> HummockVersionDeltaCommon<T>
556where
557 T: From<PbSstableInfo>,
558 PbSstableInfo: for<'a> From<&'a T>,
559{
560 pub fn from_persisted_protobuf_owned(delta: PbHummockVersionDelta) -> Self {
563 warn_if_legacy_change_log_delta_is_present(&delta);
564 delta.into()
565 }
566}
567
568fn warn_if_legacy_change_log_delta_is_present(delta: &PbHummockVersionDelta) {
569 if !delta.change_log_delta.is_empty() {
570 warn!(
571 version_delta_id = ?delta.id,
572 table_count = delta.change_log_delta.len(),
573 "deprecated table change log delta found in persisted hummock version delta; ignoring it"
574 );
575 }
576}
577
578pub trait SstableIdReader {
579 fn sst_id(&self) -> HummockSstableId;
580}
581
582pub trait ObjectIdReader {
583 fn object_id(&self) -> HummockSstableObjectId;
584}
585
586impl<T> HummockVersionDeltaCommon<T>
587where
588 T: SstableIdReader + ObjectIdReader,
589{
590 pub fn newly_added_object_ids(&self) -> HashSet<HummockObjectId> {
595 match HummockObjectId::Sstable(0.into()) {
599 HummockObjectId::Sstable(_) => {}
600 HummockObjectId::VectorFile(_) => {}
601 HummockObjectId::HnswGraphFile(_) => {}
602 };
603 self.newly_added_sst_infos()
604 .map(|sst| HummockObjectId::Sstable(sst.object_id()))
605 .chain(
606 self.vector_index_delta
607 .values()
608 .flat_map(|vector_index_delta| {
609 vector_index_delta
610 .newly_added_objects()
611 .map(|(object_id, _)| object_id)
612 }),
613 )
614 .collect()
615 }
616
617 pub fn newly_added_sst_ids(&self) -> HashSet<HummockSstableId> {
618 self.newly_added_sst_infos()
619 .map(|sst| sst.sst_id())
620 .collect()
621 }
622}
623
624impl<T> HummockVersionDeltaCommon<T> {
625 pub fn newly_added_sst_infos(&self) -> impl Iterator<Item = &'_ T> {
626 self.group_deltas.values().flat_map(|group_deltas| {
627 group_deltas.group_deltas.iter().flat_map(|group_delta| {
628 let sst_slice = match &group_delta {
629 GroupDeltaCommon::NewL0SubLevel(inserted_table_infos)
630 | GroupDeltaCommon::IntraLevel(IntraLevelDeltaCommon {
631 inserted_table_infos,
632 ..
633 }) => Some(inserted_table_infos.iter()),
634 GroupDeltaCommon::GroupConstruct(_)
635 | GroupDeltaCommon::GroupDestroy(_)
636 | GroupDeltaCommon::GroupMerge(_)
637 | GroupDeltaCommon::PruneTableIdsFromSsts(_) => None,
638 };
639 sst_slice.into_iter().flatten()
640 })
641 })
642 }
643}
644
645impl HummockVersionDelta {
646 #[expect(deprecated)]
647 pub fn max_committed_epoch_for_migration(&self) -> HummockEpoch {
648 self.max_committed_epoch
649 }
650}
651
652impl<T> From<&PbHummockVersionDelta> for HummockVersionDeltaCommon<T>
653where
654 T: for<'a> From<&'a PbSstableInfo>,
655{
656 fn from(pb_version_delta: &PbHummockVersionDelta) -> Self {
657 #[expect(deprecated)]
658 Self {
659 id: pb_version_delta.id,
660 prev_id: pb_version_delta.prev_id,
661 group_deltas: pb_version_delta
662 .group_deltas
663 .iter()
664 .map(|(group_id, deltas)| (*group_id, GroupDeltasCommon::from(deltas)))
665 .collect(),
666 max_committed_epoch: pb_version_delta.max_committed_epoch,
667 trivial_move: pb_version_delta.trivial_move,
668 new_table_watermarks: pb_version_delta
669 .new_table_watermarks
670 .iter()
671 .map(|(table_id, watermarks)| (*table_id, TableWatermarks::from(watermarks)))
672 .collect(),
673 removed_table_ids: pb_version_delta.removed_table_ids.iter().copied().collect(),
674 state_table_info_delta: pb_version_delta
675 .state_table_info_delta
676 .iter()
677 .map(|(table_id, delta)| (*table_id, *delta))
678 .collect(),
679 vector_index_delta: pb_version_delta
680 .vector_index_delta
681 .iter()
682 .map(|(table_id, delta)| (*table_id, delta.clone().into()))
683 .collect(),
684 }
685 }
686}
687
688impl<T> From<&HummockVersionDeltaCommon<T>> for PbHummockVersionDelta
689where
690 PbSstableInfo: for<'a> From<&'a T>,
691{
692 fn from(version_delta: &HummockVersionDeltaCommon<T>) -> Self {
693 #[expect(deprecated)]
694 Self {
695 id: version_delta.id,
696 prev_id: version_delta.prev_id,
697 group_deltas: version_delta
698 .group_deltas
699 .iter()
700 .map(|(group_id, deltas)| (*group_id, deltas.into()))
701 .collect(),
702 max_committed_epoch: version_delta.max_committed_epoch,
703 trivial_move: version_delta.trivial_move,
704 new_table_watermarks: version_delta
705 .new_table_watermarks
706 .iter()
707 .map(|(table_id, watermarks)| (*table_id, watermarks.into()))
708 .collect(),
709 removed_table_ids: version_delta.removed_table_ids.iter().copied().collect(),
710 change_log_delta: Default::default(),
711 state_table_info_delta: version_delta.state_table_info_delta.clone(),
712 vector_index_delta: version_delta
713 .vector_index_delta
714 .iter()
715 .map(|(table_id, delta)| (*table_id, delta.clone().into()))
716 .collect(),
717 }
718 }
719}
720
721impl<T> From<HummockVersionDeltaCommon<T>> for PbHummockVersionDelta
722where
723 PbSstableInfo: From<T>,
724{
725 fn from(version_delta: HummockVersionDeltaCommon<T>) -> Self {
726 #[expect(deprecated)]
727 Self {
728 id: version_delta.id,
729 prev_id: version_delta.prev_id,
730 group_deltas: version_delta
731 .group_deltas
732 .into_iter()
733 .map(|(group_id, deltas)| (group_id, deltas.into()))
734 .collect(),
735 max_committed_epoch: version_delta.max_committed_epoch,
736 trivial_move: version_delta.trivial_move,
737 new_table_watermarks: version_delta
738 .new_table_watermarks
739 .into_iter()
740 .map(|(table_id, watermarks)| (table_id, watermarks.into()))
741 .collect(),
742 removed_table_ids: version_delta.removed_table_ids.into_iter().collect(),
743 change_log_delta: Default::default(),
744 state_table_info_delta: version_delta.state_table_info_delta,
745 vector_index_delta: version_delta
746 .vector_index_delta
747 .into_iter()
748 .map(|(table_id, delta)| (table_id, delta.into()))
749 .collect(),
750 }
751 }
752}
753
754impl<T> From<PbHummockVersionDelta> for HummockVersionDeltaCommon<T>
755where
756 T: From<PbSstableInfo>,
757{
758 fn from(pb_version_delta: PbHummockVersionDelta) -> Self {
759 #[expect(deprecated)]
760 Self {
761 id: pb_version_delta.id,
762 prev_id: pb_version_delta.prev_id,
763 group_deltas: pb_version_delta
764 .group_deltas
765 .into_iter()
766 .map(|(group_id, deltas)| (group_id, deltas.into()))
767 .collect(),
768 max_committed_epoch: pb_version_delta.max_committed_epoch,
769 trivial_move: pb_version_delta.trivial_move,
770 new_table_watermarks: pb_version_delta
771 .new_table_watermarks
772 .into_iter()
773 .map(|(table_id, watermarks)| (table_id, watermarks.into()))
774 .collect(),
775 removed_table_ids: pb_version_delta.removed_table_ids.into_iter().collect(),
776 state_table_info_delta: pb_version_delta.state_table_info_delta,
777 vector_index_delta: pb_version_delta
778 .vector_index_delta
779 .into_iter()
780 .map(|(table_id, delta)| (table_id, delta.into()))
781 .collect(),
782 }
783 }
784}
785
786#[derive(Debug, PartialEq, Clone)]
787pub struct IntraLevelDeltaCommon<T> {
788 pub level_idx: u32,
789 pub l0_sub_level_id: u64,
790 pub removed_table_ids: HashSet<HummockSstableId>,
791 pub inserted_table_infos: Vec<T>,
792 pub vnode_partition_count: u32,
793 pub compaction_group_version_id: u64,
794}
795
796pub type IntraLevelDelta = IntraLevelDeltaCommon<SstableInfo>;
797
798impl IntraLevelDelta {
799 pub fn estimated_encode_len(&self) -> usize {
800 size_of::<u32>()
801 + size_of::<u64>()
802 + self.removed_table_ids.len() * size_of::<u32>()
803 + self
804 .inserted_table_infos
805 .iter()
806 .map(|sst| sst.estimated_encode_len())
807 .sum::<usize>()
808 + size_of::<u32>()
809 }
810}
811
812impl<T> From<PbIntraLevelDelta> for IntraLevelDeltaCommon<T>
813where
814 T: From<PbSstableInfo>,
815{
816 fn from(pb_intra_level_delta: PbIntraLevelDelta) -> Self {
817 Self {
818 level_idx: pb_intra_level_delta.level_idx,
819 l0_sub_level_id: pb_intra_level_delta.l0_sub_level_id,
820 removed_table_ids: HashSet::from_iter(
821 pb_intra_level_delta.removed_table_ids.iter().copied(),
822 ),
823 inserted_table_infos: pb_intra_level_delta
824 .inserted_table_infos
825 .into_iter()
826 .map(Into::into)
827 .collect_vec(),
828 vnode_partition_count: pb_intra_level_delta.vnode_partition_count,
829 compaction_group_version_id: pb_intra_level_delta.compaction_group_version_id,
830 }
831 }
832}
833
834impl<T> From<IntraLevelDeltaCommon<T>> for PbIntraLevelDelta
835where
836 PbSstableInfo: From<T>,
837{
838 fn from(intra_level_delta: IntraLevelDeltaCommon<T>) -> Self {
839 Self {
840 level_idx: intra_level_delta.level_idx,
841 l0_sub_level_id: intra_level_delta.l0_sub_level_id,
842 removed_table_ids: intra_level_delta.removed_table_ids.into_iter().collect(),
843 inserted_table_infos: intra_level_delta
844 .inserted_table_infos
845 .into_iter()
846 .map(Into::into)
847 .collect_vec(),
848 vnode_partition_count: intra_level_delta.vnode_partition_count,
849 compaction_group_version_id: intra_level_delta.compaction_group_version_id,
850 }
851 }
852}
853
854impl<T> From<&IntraLevelDeltaCommon<T>> for PbIntraLevelDelta
855where
856 PbSstableInfo: for<'a> From<&'a T>,
857{
858 fn from(intra_level_delta: &IntraLevelDeltaCommon<T>) -> Self {
859 Self {
860 level_idx: intra_level_delta.level_idx,
861 l0_sub_level_id: intra_level_delta.l0_sub_level_id,
862 removed_table_ids: intra_level_delta
863 .removed_table_ids
864 .iter()
865 .copied()
866 .collect(),
867 inserted_table_infos: intra_level_delta
868 .inserted_table_infos
869 .iter()
870 .map(Into::into)
871 .collect_vec(),
872 vnode_partition_count: intra_level_delta.vnode_partition_count,
873 compaction_group_version_id: intra_level_delta.compaction_group_version_id,
874 }
875 }
876}
877
878impl<T> From<&PbIntraLevelDelta> for IntraLevelDeltaCommon<T>
879where
880 T: for<'a> From<&'a PbSstableInfo>,
881{
882 fn from(pb_intra_level_delta: &PbIntraLevelDelta) -> Self {
883 Self {
884 level_idx: pb_intra_level_delta.level_idx,
885 l0_sub_level_id: pb_intra_level_delta.l0_sub_level_id,
886 removed_table_ids: HashSet::from_iter(
887 pb_intra_level_delta.removed_table_ids.iter().copied(),
888 ),
889 inserted_table_infos: pb_intra_level_delta
890 .inserted_table_infos
891 .iter()
892 .map(Into::into)
893 .collect_vec(),
894 vnode_partition_count: pb_intra_level_delta.vnode_partition_count,
895 compaction_group_version_id: pb_intra_level_delta.compaction_group_version_id,
896 }
897 }
898}
899
900impl IntraLevelDelta {
901 pub fn new(
902 level_idx: u32,
903 l0_sub_level_id: u64,
904 removed_table_ids: HashSet<HummockSstableId>,
905 inserted_table_infos: Vec<SstableInfo>,
906 vnode_partition_count: u32,
907 compaction_group_version_id: u64,
908 ) -> Self {
909 Self {
910 level_idx,
911 l0_sub_level_id,
912 removed_table_ids,
913 inserted_table_infos,
914 vnode_partition_count,
915 compaction_group_version_id,
916 }
917 }
918}
919
920#[derive(Debug, PartialEq, Clone)]
921pub enum GroupDeltaCommon<T> {
922 NewL0SubLevel(Vec<T>),
923 IntraLevel(IntraLevelDeltaCommon<T>),
924 GroupConstruct(Box<PbGroupConstruct>),
925 GroupDestroy(PbGroupDestroy),
926 GroupMerge(PbGroupMerge),
927 PruneTableIdsFromSsts(HashSet<TableId>),
931}
932
933pub type GroupDelta = GroupDeltaCommon<SstableInfo>;
934
935impl<T> From<PbGroupDelta> for GroupDeltaCommon<T>
936where
937 T: From<PbSstableInfo>,
938{
939 fn from(pb_group_delta: PbGroupDelta) -> Self {
940 match pb_group_delta.delta_type {
941 Some(PbDeltaType::IntraLevel(pb_intra_level_delta)) => {
942 GroupDeltaCommon::IntraLevel(IntraLevelDeltaCommon::from(pb_intra_level_delta))
943 }
944 Some(PbDeltaType::GroupConstruct(pb_group_construct)) => {
945 GroupDeltaCommon::GroupConstruct(Box::new(pb_group_construct))
946 }
947 Some(PbDeltaType::GroupDestroy(pb_group_destroy)) => {
948 GroupDeltaCommon::GroupDestroy(pb_group_destroy)
949 }
950 Some(PbDeltaType::GroupMerge(pb_group_merge)) => {
951 GroupDeltaCommon::GroupMerge(pb_group_merge)
952 }
953 Some(DeltaType::NewL0SubLevel(pb_new_sub_level)) => GroupDeltaCommon::NewL0SubLevel(
954 pb_new_sub_level
955 .inserted_table_infos
956 .into_iter()
957 .map(T::from)
958 .collect(),
959 ),
960 Some(PbDeltaType::TruncateTables(pb_truncate_tables)) => {
961 GroupDeltaCommon::PruneTableIdsFromSsts(
962 pb_truncate_tables.table_ids.into_iter().collect(),
963 )
964 }
965
966 None => panic!("delta_type is not set"),
967 }
968 }
969}
970
971impl<T> From<GroupDeltaCommon<T>> for PbGroupDelta
972where
973 PbSstableInfo: From<T>,
974{
975 fn from(group_delta: GroupDeltaCommon<T>) -> Self {
976 match group_delta {
977 GroupDeltaCommon::IntraLevel(intra_level_delta) => PbGroupDelta {
978 delta_type: Some(PbDeltaType::IntraLevel(intra_level_delta.into())),
979 },
980 GroupDeltaCommon::GroupConstruct(pb_group_construct) => PbGroupDelta {
981 delta_type: Some(PbDeltaType::GroupConstruct(*pb_group_construct)),
982 },
983 GroupDeltaCommon::GroupDestroy(pb_group_destroy) => PbGroupDelta {
984 delta_type: Some(PbDeltaType::GroupDestroy(pb_group_destroy)),
985 },
986 GroupDeltaCommon::GroupMerge(pb_group_merge) => PbGroupDelta {
987 delta_type: Some(PbDeltaType::GroupMerge(pb_group_merge)),
988 },
989 GroupDeltaCommon::NewL0SubLevel(new_sub_level) => PbGroupDelta {
990 delta_type: Some(PbDeltaType::NewL0SubLevel(PbNewL0SubLevel {
991 inserted_table_infos: new_sub_level
992 .into_iter()
993 .map(PbSstableInfo::from)
994 .collect(),
995 })),
996 },
997 GroupDeltaCommon::PruneTableIdsFromSsts(table_ids) => PbGroupDelta {
998 delta_type: Some(PbDeltaType::TruncateTables(PbTruncateTables {
999 table_ids: table_ids.iter().copied().collect(),
1000 })),
1001 },
1002 }
1003 }
1004}
1005
1006impl<T> From<&GroupDeltaCommon<T>> for PbGroupDelta
1007where
1008 PbSstableInfo: for<'a> From<&'a T>,
1009{
1010 fn from(group_delta: &GroupDeltaCommon<T>) -> Self {
1011 match group_delta {
1012 GroupDeltaCommon::IntraLevel(intra_level_delta) => PbGroupDelta {
1013 delta_type: Some(PbDeltaType::IntraLevel(intra_level_delta.into())),
1014 },
1015 GroupDeltaCommon::GroupConstruct(pb_group_construct) => PbGroupDelta {
1016 delta_type: Some(PbDeltaType::GroupConstruct(*pb_group_construct.clone())),
1017 },
1018 GroupDeltaCommon::GroupDestroy(pb_group_destroy) => PbGroupDelta {
1019 delta_type: Some(PbDeltaType::GroupDestroy(*pb_group_destroy)),
1020 },
1021 GroupDeltaCommon::GroupMerge(pb_group_merge) => PbGroupDelta {
1022 delta_type: Some(PbDeltaType::GroupMerge(*pb_group_merge)),
1023 },
1024 GroupDeltaCommon::NewL0SubLevel(new_sub_level) => PbGroupDelta {
1025 delta_type: Some(PbDeltaType::NewL0SubLevel(PbNewL0SubLevel {
1026 inserted_table_infos: new_sub_level.iter().map(PbSstableInfo::from).collect(),
1027 })),
1028 },
1029 GroupDeltaCommon::PruneTableIdsFromSsts(table_ids) => PbGroupDelta {
1030 delta_type: Some(PbDeltaType::TruncateTables(PbTruncateTables {
1031 table_ids: table_ids.iter().copied().collect(),
1032 })),
1033 },
1034 }
1035 }
1036}
1037
1038impl<T> From<&PbGroupDelta> for GroupDeltaCommon<T>
1039where
1040 T: for<'a> From<&'a PbSstableInfo>,
1041{
1042 fn from(pb_group_delta: &PbGroupDelta) -> Self {
1043 match &pb_group_delta.delta_type {
1044 Some(PbDeltaType::IntraLevel(pb_intra_level_delta)) => {
1045 GroupDeltaCommon::IntraLevel(IntraLevelDeltaCommon::from(pb_intra_level_delta))
1046 }
1047 Some(PbDeltaType::GroupConstruct(pb_group_construct)) => {
1048 GroupDeltaCommon::GroupConstruct(Box::new(pb_group_construct.clone()))
1049 }
1050 Some(PbDeltaType::GroupDestroy(pb_group_destroy)) => {
1051 GroupDeltaCommon::GroupDestroy(*pb_group_destroy)
1052 }
1053 Some(PbDeltaType::GroupMerge(pb_group_merge)) => {
1054 GroupDeltaCommon::GroupMerge(*pb_group_merge)
1055 }
1056 Some(DeltaType::NewL0SubLevel(pb_new_sub_level)) => GroupDeltaCommon::NewL0SubLevel(
1057 pb_new_sub_level
1058 .inserted_table_infos
1059 .iter()
1060 .map(T::from)
1061 .collect(),
1062 ),
1063 Some(PbDeltaType::TruncateTables(pb_truncate_tables)) => {
1064 GroupDeltaCommon::PruneTableIdsFromSsts(
1065 pb_truncate_tables.table_ids.iter().copied().collect(),
1066 )
1067 }
1068 None => panic!("delta_type is not set"),
1069 }
1070 }
1071}
1072
1073#[derive(Debug, PartialEq, Clone)]
1074pub struct GroupDeltasCommon<T> {
1075 pub group_deltas: Vec<GroupDeltaCommon<T>>,
1076}
1077
1078impl<T> Default for GroupDeltasCommon<T> {
1079 fn default() -> Self {
1080 Self {
1081 group_deltas: vec![],
1082 }
1083 }
1084}
1085
1086pub type GroupDeltas = GroupDeltasCommon<SstableInfo>;
1087
1088impl<T> From<PbGroupDeltas> for GroupDeltasCommon<T>
1089where
1090 T: From<PbSstableInfo>,
1091{
1092 fn from(pb_group_deltas: PbGroupDeltas) -> Self {
1093 Self {
1094 group_deltas: pb_group_deltas
1095 .group_deltas
1096 .into_iter()
1097 .map(GroupDeltaCommon::from)
1098 .collect_vec(),
1099 }
1100 }
1101}
1102
1103impl<T> From<GroupDeltasCommon<T>> for PbGroupDeltas
1104where
1105 PbSstableInfo: From<T>,
1106{
1107 fn from(group_deltas: GroupDeltasCommon<T>) -> Self {
1108 Self {
1109 group_deltas: group_deltas
1110 .group_deltas
1111 .into_iter()
1112 .map(|group_delta| group_delta.into())
1113 .collect_vec(),
1114 }
1115 }
1116}
1117
1118impl<T> From<&GroupDeltasCommon<T>> for PbGroupDeltas
1119where
1120 PbSstableInfo: for<'a> From<&'a T>,
1121{
1122 fn from(group_deltas: &GroupDeltasCommon<T>) -> Self {
1123 Self {
1124 group_deltas: group_deltas
1125 .group_deltas
1126 .iter()
1127 .map(|group_delta| group_delta.into())
1128 .collect_vec(),
1129 }
1130 }
1131}
1132
1133impl<T> From<&PbGroupDeltas> for GroupDeltasCommon<T>
1134where
1135 T: for<'a> From<&'a PbSstableInfo>,
1136{
1137 fn from(pb_group_deltas: &PbGroupDeltas) -> Self {
1138 Self {
1139 group_deltas: pb_group_deltas
1140 .group_deltas
1141 .iter()
1142 .map(GroupDeltaCommon::from)
1143 .collect_vec(),
1144 }
1145 }
1146}
1147
1148impl<T> GroupDeltasCommon<T>
1149where
1150 PbSstableInfo: for<'a> From<&'a T>,
1151{
1152 pub fn to_protobuf(&self) -> PbGroupDeltas {
1153 self.into()
1154 }
1155}
1156
1157#[cfg(test)]
1158mod tests {
1159 use super::*;
1160
1161 #[test]
1162 fn deprecated_change_log_delta_is_ignored() {
1163 let mut pb_delta = PbHummockVersionDelta {
1164 id: 1.into(),
1165 ..Default::default()
1166 };
1167 pb_delta
1168 .change_log_delta
1169 .insert(1.into(), Default::default());
1170
1171 let delta = HummockVersionDelta::from_persisted_protobuf_owned(pb_delta);
1172 assert!(delta.to_protobuf().change_log_delta.is_empty());
1173 }
1174}