risingwave_hummock_sdk/
time_travel.rs1use std::collections::{HashMap, HashSet};
16
17use risingwave_common::catalog::TableId;
18use risingwave_pb::hummock::PbSstableInfo;
19use risingwave_pb::hummock::hummock_version::PbLevels;
20use risingwave_pb::hummock::hummock_version_delta::PbGroupDeltas;
21
22use crate::compaction_group::StateTableId;
23use crate::level::{Level, Levels, LevelsCommon};
24use crate::sstable_info::SstableInfo;
25use crate::version::{
26 GroupDelta, GroupDeltas, GroupDeltasCommon, HummockVersion, HummockVersionCommon,
27 HummockVersionDelta, HummockVersionDeltaCommon, ObjectIdReader, SstableIdReader,
28};
29use crate::{CompactionGroupId, HummockSstableId, HummockSstableObjectId};
30
31pub type IncompleteHummockVersion = HummockVersionCommon<SstableIdInVersion>;
32
33pub fn refill_version(
36 version: &mut HummockVersion,
37 sst_id_to_info: &HashMap<HummockSstableId, SstableInfo>,
38 table_id: TableId,
39) {
40 for level in version.levels.values_mut().flat_map(|level| {
41 level
42 .l0
43 .sub_levels
44 .iter_mut()
45 .rev()
46 .chain(level.levels.iter_mut())
47 }) {
48 refill_level(level, sst_id_to_info);
49 level
50 .table_infos
51 .retain(|t| t.table_ids.contains(&table_id));
52 }
53}
54
55fn refill_level(level: &mut Level, sst_id_to_info: &HashMap<HummockSstableId, SstableInfo>) {
56 for s in &mut level.table_infos {
57 refill_sstable_info(s, sst_id_to_info);
58 }
59}
60
61fn refill_sstable_info(
63 sstable_info: &mut SstableInfo,
64 sst_id_to_info: &HashMap<HummockSstableId, SstableInfo>,
65) {
66 *sstable_info = sst_id_to_info
67 .get(&sstable_info.sst_id)
68 .unwrap_or_else(|| panic!("SstableInfo should exist"))
69 .clone();
70}
71
72impl From<(&HummockVersion, &HashSet<StateTableId>)> for IncompleteHummockVersion {
74 fn from(p: (&HummockVersion, &HashSet<StateTableId>)) -> Self {
75 let (version, time_travel_table_ids) = p;
76 #[expect(deprecated)]
77 Self {
78 id: version.id,
79 levels: version
80 .levels
81 .iter()
82 .map(|(group_id, levels)| {
83 let levels = rewrite_levels(levels, time_travel_table_ids);
84 (*group_id as CompactionGroupId, levels)
85 })
86 .collect(),
87 max_committed_epoch: version.max_committed_epoch,
88 table_watermarks: version.table_watermarks.clone(),
89 state_table_info: version.state_table_info.clone(),
90 vector_indexes: version.vector_indexes.clone(),
91 }
92 }
93}
94
95fn rewrite_levels(
97 levels: &Levels,
98 time_travel_table_ids: &HashSet<StateTableId>,
99) -> LevelsCommon<SstableIdInVersion> {
100 fn rewrite_level(level: &mut Level, time_travel_table_ids: &HashSet<StateTableId>) {
101 level.table_infos.retain(|sst| {
103 sst.table_ids
104 .iter()
105 .any(|tid| time_travel_table_ids.contains(tid))
106 });
107 }
108 let mut levels = levels.clone();
109 for level in &mut levels.levels {
110 rewrite_level(level, time_travel_table_ids);
111 }
112 {
113 let l0 = &mut levels.l0;
114 for sub_level in &mut l0.sub_levels {
115 rewrite_level(sub_level, time_travel_table_ids);
116 }
117 l0.sub_levels.retain(|s| !s.table_infos.is_empty());
118 }
119 PbLevels::from(levels).into()
120}
121
122pub type IncompleteHummockVersionDelta = HummockVersionDeltaCommon<SstableIdInVersion>;
125
126impl From<(&HummockVersionDelta, &HashSet<StateTableId>)> for IncompleteHummockVersionDelta {
128 fn from(p: (&HummockVersionDelta, &HashSet<StateTableId>)) -> Self {
129 let (delta, time_travel_table_ids) = p;
130 #[expect(deprecated)]
131 Self {
132 id: delta.id,
133 prev_id: delta.prev_id,
134 group_deltas: delta
135 .group_deltas
136 .iter()
137 .map(|(cg_id, deltas)| {
138 let deltas = rewrite_group_deltas(deltas, time_travel_table_ids);
139 (*cg_id, deltas)
140 })
141 .collect(),
142 max_committed_epoch: delta.max_committed_epoch,
143 trivial_move: delta.trivial_move,
144 new_table_watermarks: delta.new_table_watermarks.clone(),
145 removed_table_ids: delta.removed_table_ids.clone(),
146 state_table_info_delta: delta.state_table_info_delta.clone(),
147 vector_index_delta: delta.vector_index_delta.clone(),
148 }
149 }
150}
151
152fn rewrite_group_deltas(
154 group_deltas: &GroupDeltas,
155 time_travel_table_ids: &HashSet<StateTableId>,
156) -> GroupDeltasCommon<SstableIdInVersion> {
157 let mut group_deltas = group_deltas.clone();
158 for group_delta in &mut group_deltas.group_deltas {
159 let GroupDelta::NewL0SubLevel(inserted_table_infos) = group_delta else {
160 tracing::error!(?group_delta, "unexpected delta type");
161 continue;
162 };
163 inserted_table_infos.retain(|sst| {
164 sst.table_ids
165 .iter()
166 .any(|tid| time_travel_table_ids.contains(tid))
167 });
168 }
169 PbGroupDeltas::from(group_deltas).into()
170}
171
172pub struct SstableIdInVersion {
173 sst_id: HummockSstableId,
174 object_id: HummockSstableObjectId,
175}
176
177impl SstableIdReader for SstableIdInVersion {
178 fn sst_id(&self) -> HummockSstableId {
179 self.sst_id
180 }
181}
182
183impl ObjectIdReader for SstableIdInVersion {
184 fn object_id(&self) -> HummockSstableObjectId {
185 self.object_id
186 }
187}
188
189impl From<&SstableIdInVersion> for PbSstableInfo {
190 fn from(sst_id: &SstableIdInVersion) -> Self {
191 Self {
192 sst_id: sst_id.sst_id,
193 object_id: sst_id.object_id,
194 ..Default::default()
195 }
196 }
197}
198
199impl From<SstableIdInVersion> for PbSstableInfo {
200 fn from(sst_id: SstableIdInVersion) -> Self {
201 (&sst_id).into()
202 }
203}
204
205impl From<&PbSstableInfo> for SstableIdInVersion {
206 fn from(s: &PbSstableInfo) -> Self {
207 SstableIdInVersion {
208 sst_id: s.sst_id,
209 object_id: s.object_id,
210 }
211 }
212}
213
214impl From<PbSstableInfo> for SstableIdInVersion {
215 fn from(value: PbSstableInfo) -> Self {
216 (&value).into()
217 }
218}