Skip to main content

risingwave_hummock_sdk/
version.rs

1// Copyright 2023 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::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    // in memory index
47    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                        // table moved to another compaction group
163                        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    /// Convert the `PbHummockVersion` received from rpc to `HummockVersion`. No need to
249    /// maintain backward compatibility.
250    pub fn from_rpc_protobuf(pb_version: &PbHummockVersion) -> Self {
251        pb_version.into()
252    }
253
254    /// Convert the `PbHummockVersion` deserialized from persisted state to `HummockVersion`.
255    /// We should maintain backward compatibility.
256    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    /// Convert an owned `PbHummockVersion` deserialized from persisted state to `HummockVersion`,
271    /// moving data instead of cloning for better performance on large checkpoints.
272    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        // for backward-compatibility of previous hummock version delta
428        self.state_table_info.state_table_info.is_empty()
429            && self.levels.values().any(|group| {
430                // state_table_info is not previously filled, but there previously exists some tables
431                #[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 backward-compatibility of previous hummock version delta
442        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    /// Convert the `PbHummockVersionDelta` deserialized from persisted state to `HummockVersionDelta`.
538    /// We should maintain backward compatibility.
539    pub fn from_persisted_protobuf(delta: &PbHummockVersionDelta) -> Self {
540        warn_if_legacy_change_log_delta_is_present(delta);
541        delta.into()
542    }
543
544    /// Convert the `PbHummockVersionDelta` received from rpc to `HummockVersionDelta`. No need to
545    /// maintain backward compatibility.
546    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    /// Convert an owned `PbHummockVersionDelta` deserialized from persisted state to
561    /// `HummockVersionDelta`, moving data instead of cloning.
562    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    /// Get the newly added object ids from the version delta.
591    ///
592    /// Note: the result can be false positive because we only collect the set of sst object ids in the `inserted_table_infos`,
593    /// but it is possible that the object is moved or split from other compaction groups or levels.
594    pub fn newly_added_object_ids(&self) -> HashSet<HummockObjectId> {
595        // DO NOT REMOVE THIS LINE
596        // This is to ensure that when adding new variant to `HummockObjectId`,
597        // the compiler will warn us if we forget to handle it here.
598        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    /// Prunes table ids from SST metadata.
928    ///
929    /// Serialized as protobuf `truncate_tables` for compatibility.
930    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}