1use std::collections::{HashMap, HashSet, VecDeque};
16
17use anyhow::anyhow;
18use futures::TryStreamExt;
19use risingwave_common::catalog::TableId;
20use risingwave_common::system_param::reader::SystemParamsRead;
21use risingwave_common::util::epoch::Epoch;
22use risingwave_hummock_sdk::compaction_group::StateTableId;
23use risingwave_hummock_sdk::sstable_info::SstableInfo;
24use risingwave_hummock_sdk::time_travel::{
25 IncompleteHummockVersion, IncompleteHummockVersionDelta, refill_version,
26};
27use risingwave_hummock_sdk::version::{GroupDeltaCommon, HummockVersion, HummockVersionDelta};
28use risingwave_hummock_sdk::{CompactionGroupId, HummockEpoch, HummockObjectId, HummockSstableId};
29use risingwave_meta_model::hummock_sstable_info::SstableInfoV2Backend;
30use risingwave_meta_model::{
31 HummockVersionId, hummock_epoch_to_version, hummock_sstable_info, hummock_time_travel_delta,
32 hummock_time_travel_version,
33};
34use risingwave_pb::hummock::{PbHummockVersion, PbHummockVersionDelta};
35use sea_orm::ActiveValue::Set;
36use sea_orm::{
37 ColumnTrait, Condition, ConnectionTrait, DatabaseTransaction, EntityTrait, PaginatorTrait,
38 QueryFilter, QueryOrder, QuerySelect, TransactionTrait,
39};
40use tracing::info;
41
42use crate::hummock::HummockManager;
43use crate::hummock::error::{Error, Result};
44
45impl HummockManager {
47 pub(crate) async fn init_time_travel_state(&self) -> Result<()> {
48 let sql_store = self.env.meta_store_ref();
49 let mut guard = self.versioning.write().await;
50 guard.mark_next_time_travel_version_snapshot();
51
52 guard.last_time_travel_snapshot_sst_ids = HashSet::new();
53 let Some(version) = hummock_time_travel_version::Entity::find()
54 .order_by_desc(hummock_time_travel_version::Column::VersionId)
55 .one(&sql_store.conn)
56 .await?
57 .map(|v| {
58 IncompleteHummockVersion::from_persisted_protobuf_owned(v.version.to_protobuf())
59 })
60 else {
61 return Ok(());
62 };
63 guard.last_time_travel_snapshot_sst_ids = version.get_sst_ids();
64 Ok(())
65 }
66
67 pub(crate) async fn truncate_time_travel_metadata(
68 &self,
69 epoch_watermark: HummockEpoch,
70 pinned_snapshot_epochs: HashMap<TableId, HashSet<HummockEpoch>>,
71 ) -> Result<()> {
72 let _timer = self
73 .metrics
74 .time_travel_vacuum_metadata_latency
75 .start_timer();
76 let min_pinned_version_id = self.context_info.read().await.min_pinned_version_id();
77 let sql_store = self.env.meta_store_ref();
78 let txn = sql_store.conn.begin().await?;
79
80 let epoch_watermark_model =
81 risingwave_meta_model::Epoch::try_from(epoch_watermark).unwrap();
82 let has_expired_epoch = hummock_epoch_to_version::Entity::find()
83 .filter(hummock_epoch_to_version::Column::Epoch.lt(epoch_watermark_model))
84 .select_only()
85 .column(hummock_epoch_to_version::Column::VersionId)
86 .into_tuple::<HummockVersionId>()
87 .one(&txn)
88 .await?
89 .is_some();
90 if !has_expired_epoch {
91 txn.commit().await?;
92 return Ok(());
93 }
94
95 let mut pinned_epoch_rows = HashSet::new();
96 let mut pinned_snapshot_version_ids = HashSet::new();
97 for (table_id, pinned_epochs) in pinned_snapshot_epochs {
98 for pinned_epoch in pinned_epochs {
99 if pinned_epoch >= epoch_watermark {
100 continue;
101 }
102 let pinned_epoch_model =
103 risingwave_meta_model::Epoch::try_from(pinned_epoch).unwrap();
104 let epoch_to_version = hummock_epoch_to_version::Entity::find_by_id((
105 pinned_epoch_model,
106 i64::from(table_id.as_raw_id()),
107 ))
108 .one(&txn)
109 .await?;
110 let Some(epoch_to_version) = epoch_to_version else {
111 tracing::warn!(
112 %table_id,
113 pinned_epoch,
114 "pinned snapshot epoch mapping not found, skip pinning"
115 );
116 continue;
117 };
118 pinned_epoch_rows.insert((epoch_to_version.epoch, epoch_to_version.table_id));
119 pinned_snapshot_version_ids.insert(epoch_to_version.version_id);
120 }
121 }
122
123 let mut pinned_replay_version_ids = HashSet::new();
124 let mut pinned_delta_ids = HashSet::new();
125 let mut pinned_snapshot_sst_ids = HashSet::new();
126 let mut pinned_snapshot_object_ids = HashSet::new();
127 for pinned_snapshot_version_id in pinned_snapshot_version_ids {
128 let Some(resolved) =
129 resolve_time_travel_version(&txn, pinned_snapshot_version_id).await?
130 else {
131 tracing::warn!(
132 %pinned_snapshot_version_id,
133 "time travel version before pinned version not found, skip pinning"
134 );
135 continue;
136 };
137 pinned_replay_version_ids.insert(resolved.replay_version.id);
138 pinned_snapshot_sst_ids.extend(resolved.replay_version.get_sst_ids());
139 pinned_snapshot_object_ids.extend(resolved.replay_version.get_object_ids());
140 for delta in resolved.deltas {
141 pinned_delta_ids.insert(delta.id);
142 pinned_snapshot_sst_ids.extend(delta.newly_added_sst_ids(true));
143 pinned_snapshot_object_ids.extend(delta.newly_added_object_ids(true));
144 }
145 }
146
147 let version_watermark = hummock_epoch_to_version::Entity::find()
152 .filter(hummock_epoch_to_version::Column::Epoch.gte(epoch_watermark_model))
153 .select_only()
154 .column(hummock_epoch_to_version::Column::VersionId)
155 .order_by_asc(hummock_epoch_to_version::Column::VersionId)
156 .into_tuple::<HummockVersionId>()
157 .one(&txn)
158 .await?;
159 let mut watermark_version_id = version_watermark.map_or(min_pinned_version_id, |id| {
161 std::cmp::min(id, min_pinned_version_id)
162 });
163 if let Some(max_version_count) = self.env.opts.time_travel_vacuum_max_version_count {
164 let mut query = hummock_time_travel_version::Entity::find()
165 .select_only()
166 .column(hummock_time_travel_version::Column::VersionId)
167 .order_by_asc(hummock_time_travel_version::Column::VersionId)
168 .limit(2);
169 if !pinned_replay_version_ids.is_empty() {
172 query = query.filter(
173 hummock_time_travel_version::Column::VersionId
174 .is_not_in(pinned_replay_version_ids.iter().copied()),
175 );
176 }
177 let earliest2_version_ids = query.into_tuple::<HummockVersionId>().all(&txn).await?;
178 if earliest2_version_ids.len() == 2 {
180 watermark_version_id = std::cmp::min(
181 watermark_version_id,
182 HummockVersionId::new(std::cmp::max(
183 earliest2_version_ids[0]
184 .as_raw_id()
185 .saturating_add(max_version_count.into()),
186 earliest2_version_ids[1].as_raw_id(),
187 )),
188 );
189 }
190 }
191 let mut delete_epoch_rows = hummock_epoch_to_version::Entity::delete_many()
192 .filter(hummock_epoch_to_version::Column::Epoch.lt(epoch_watermark_model));
193 for (epoch, table_id) in &pinned_epoch_rows {
194 delete_epoch_rows = delete_epoch_rows.filter(
195 Condition::any()
196 .add(hummock_epoch_to_version::Column::Epoch.ne(*epoch))
197 .add(hummock_epoch_to_version::Column::TableId.ne(*table_id)),
198 );
199 }
200 let res = delete_epoch_rows.exec(&txn).await?;
201 tracing::info!(
202 epoch_watermark,
203 "Delete {} rows from hummock_epoch_to_version.",
204 res.rows_affected
205 );
206 let latest_valid_version = hummock_time_travel_version::Entity::find()
207 .filter(hummock_time_travel_version::Column::VersionId.lte(watermark_version_id))
208 .order_by_desc(hummock_time_travel_version::Column::VersionId)
209 .one(&txn)
210 .await?
211 .map(|m| {
212 IncompleteHummockVersion::from_persisted_protobuf_owned(m.version.to_protobuf())
213 });
214 let Some(latest_valid_version) = latest_valid_version else {
215 txn.commit().await?;
216 return Ok(());
217 };
218 let latest_valid_version_id = latest_valid_version.id;
219 let mut retained_snapshot_sst_ids = latest_valid_version.get_sst_ids();
220 let mut retained_snapshot_object_ids = latest_valid_version
221 .get_object_ids()
222 .collect::<HashSet<_>>();
223 retained_snapshot_sst_ids.extend(pinned_snapshot_sst_ids);
224 retained_snapshot_object_ids.extend(pinned_snapshot_object_ids);
225 let mut object_ids_to_delete: HashSet<_> = HashSet::default();
226 let mut version_delete_condition = Condition::all()
227 .add(hummock_time_travel_version::Column::VersionId.lt(latest_valid_version_id));
228 if !pinned_replay_version_ids.is_empty() {
229 version_delete_condition = version_delete_condition.add(
230 hummock_time_travel_version::Column::VersionId
231 .is_not_in(pinned_replay_version_ids.iter().copied()),
232 );
233 }
234 let version_ids_to_delete: Vec<risingwave_meta_model::HummockVersionId> =
235 hummock_time_travel_version::Entity::find()
236 .select_only()
237 .column(hummock_time_travel_version::Column::VersionId)
238 .filter(version_delete_condition.clone())
239 .order_by_desc(hummock_time_travel_version::Column::VersionId)
240 .into_tuple()
241 .all(&txn)
242 .await?;
243 let mut delta_delete_condition = Condition::all()
244 .add(hummock_time_travel_delta::Column::VersionId.lt(latest_valid_version_id));
245 if !pinned_delta_ids.is_empty() {
246 delta_delete_condition = delta_delete_condition.add(
247 hummock_time_travel_delta::Column::VersionId
248 .is_not_in(pinned_delta_ids.iter().copied()),
249 );
250 }
251 let delta_ids_to_delete: Vec<risingwave_meta_model::HummockVersionId> =
252 hummock_time_travel_delta::Entity::find()
253 .select_only()
254 .column(hummock_time_travel_delta::Column::VersionId)
255 .filter(delta_delete_condition.clone())
256 .into_tuple()
257 .all(&txn)
258 .await?;
259 let delete_sst_batch_size = self
260 .env
261 .opts
262 .hummock_time_travel_epoch_version_insert_batch_size;
263 let delta_fetch_batch_size = self
264 .env
265 .opts
266 .hummock_time_travel_delta_fetch_batch_size
267 .max(1);
268 let mut sst_ids_to_delete: HashSet<_> = HashSet::default();
269 async fn delete_sst_in_batch(
270 txn: &DatabaseTransaction,
271 sst_ids_to_delete: HashSet<HummockSstableId>,
272 delete_sst_batch_size: usize,
273 ) -> Result<()> {
274 for start_idx in 0..=(sst_ids_to_delete.len().saturating_sub(1) / delete_sst_batch_size)
275 {
276 hummock_sstable_info::Entity::delete_many()
277 .filter(
278 hummock_sstable_info::Column::SstId.is_in(
279 sst_ids_to_delete
280 .iter()
281 .skip(start_idx * delete_sst_batch_size)
282 .take(delete_sst_batch_size)
283 .copied(),
284 ),
285 )
286 .exec(txn)
287 .await?;
288 }
289 Ok(())
290 }
291 for delta_id_batch in delta_ids_to_delete.chunks(delta_fetch_batch_size) {
292 let mut delta_to_delete_by_id: HashMap<_, _> =
293 hummock_time_travel_delta::Entity::find()
294 .filter(
295 hummock_time_travel_delta::Column::VersionId
296 .is_in(delta_id_batch.iter().copied()),
297 )
298 .all(&txn)
299 .await?
300 .into_iter()
301 .map(|delta| (delta.version_id, delta))
302 .collect();
303 for &delta_id_to_delete in delta_id_batch {
304 let delta_to_delete = delta_to_delete_by_id
305 .remove(&delta_id_to_delete)
306 .ok_or_else(|| {
307 Error::TimeTravel(anyhow!(format!(
308 "version delta {} not found",
309 delta_id_to_delete
310 )))
311 })?;
312 let delta_to_delete = IncompleteHummockVersionDelta::from_persisted_protobuf_owned(
313 delta_to_delete.version_delta.to_protobuf(),
314 );
315 let new_sst_ids = delta_to_delete.newly_added_sst_ids(true);
316 sst_ids_to_delete.extend(&new_sst_ids - &retained_snapshot_sst_ids);
318 if sst_ids_to_delete.len() >= delete_sst_batch_size {
319 delete_sst_in_batch(
320 &txn,
321 std::mem::take(&mut sst_ids_to_delete),
322 delete_sst_batch_size,
323 )
324 .await?;
325 }
326 let new_object_ids = delta_to_delete.newly_added_object_ids(true);
327 object_ids_to_delete.extend(&new_object_ids - &retained_snapshot_object_ids);
328 }
329 }
330 for prev_version_id in version_ids_to_delete {
331 let prev_version = {
332 let prev_version = hummock_time_travel_version::Entity::find_by_id(prev_version_id)
333 .one(&txn)
334 .await?
335 .ok_or_else(|| {
336 Error::TimeTravel(anyhow!(format!(
337 "prev_version {} not found",
338 prev_version_id
339 )))
340 })?;
341 IncompleteHummockVersion::from_persisted_protobuf_owned(
342 prev_version.version.to_protobuf(),
343 )
344 };
345 let sst_ids = prev_version.get_sst_ids();
346 sst_ids_to_delete.extend(&sst_ids - &retained_snapshot_sst_ids);
347 if sst_ids_to_delete.len() >= delete_sst_batch_size {
348 delete_sst_in_batch(
349 &txn,
350 std::mem::take(&mut sst_ids_to_delete),
351 delete_sst_batch_size,
352 )
353 .await?;
354 }
355 let new_object_ids: HashSet<_> = prev_version.get_object_ids().collect();
356 object_ids_to_delete.extend(&new_object_ids - &retained_snapshot_object_ids);
357 }
358 if !sst_ids_to_delete.is_empty() {
359 delete_sst_in_batch(&txn, sst_ids_to_delete, delete_sst_batch_size).await?;
360 }
361
362 if !object_ids_to_delete.is_empty() {
363 self.gc_manager
366 .add_may_delete_object_ids(object_ids_to_delete.into_iter());
367 }
368
369 let res = hummock_time_travel_version::Entity::delete_many()
370 .filter(version_delete_condition)
371 .exec(&txn)
372 .await?;
373 tracing::info!(
374 %watermark_version_id,
375 %latest_valid_version_id,
376 "Deleted {} rows from hummock_time_travel_version.",
377 res.rows_affected
378 );
379
380 let res = hummock_time_travel_delta::Entity::delete_many()
381 .filter(delta_delete_condition)
382 .exec(&txn)
383 .await?;
384 tracing::info!(
385 %watermark_version_id,
386 %latest_valid_version_id,
387 "Deleted {} rows from hummock_time_travel_delta.",
388 res.rows_affected
389 );
390
391 txn.commit().await?;
392 Ok(())
393 }
394
395 pub(crate) async fn filter_out_objects_by_time_travel_v1(
396 &self,
397 objects: impl Iterator<Item = HummockObjectId>,
398 ) -> Result<HashSet<HummockObjectId>> {
399 let batch_size = self
400 .env
401 .opts
402 .hummock_time_travel_filter_out_objects_batch_size;
403 info!("filter out objects by time travel v1, only sst will remain in the result set");
404 let mut result: HashSet<_> = objects
407 .filter(|object_id| match object_id {
408 HummockObjectId::Sstable(_) => true,
409 HummockObjectId::VectorFile(_) | HummockObjectId::HnswGraphFile(_) => false,
410 })
411 .collect();
412 let mut remain_sst: VecDeque<_> = result.iter().copied().collect();
413 while !remain_sst.is_empty() {
414 let batch = remain_sst
415 .drain(..std::cmp::min(remain_sst.len(), batch_size))
416 .map(|object_id| object_id.as_raw());
417 let reject_object_ids: Vec<risingwave_meta_model::HummockSstableObjectId> =
418 hummock_sstable_info::Entity::find()
419 .filter(hummock_sstable_info::Column::ObjectId.is_in(batch))
420 .select_only()
421 .column(hummock_sstable_info::Column::ObjectId)
422 .into_tuple()
423 .all(&self.env.meta_store_ref().conn)
424 .await?;
425 for reject in reject_object_ids {
426 let object_id = HummockObjectId::Sstable(reject);
427 result.remove(&object_id);
428 }
429 }
430 Ok(result)
431 }
432
433 pub(crate) async fn filter_out_objects_by_time_travel(
434 &self,
435 objects: impl Iterator<Item = HummockObjectId>,
436 ) -> Result<HashSet<HummockObjectId>> {
437 if self.env.opts.hummock_time_travel_filter_out_objects_v1 {
438 return self.filter_out_objects_by_time_travel_v1(objects).await;
439 }
440 let mut result: HashSet<_> = objects.collect();
441
442 {
444 let mut prev_version_id: Option<HummockVersionId> = None;
445 loop {
446 let query = hummock_time_travel_version::Entity::find();
447 let query = if let Some(prev_version_id) = prev_version_id {
448 query.filter(hummock_time_travel_version::Column::VersionId.gt(prev_version_id))
449 } else {
450 query
451 };
452 let mut version_stream = query
453 .order_by_asc(hummock_time_travel_version::Column::VersionId)
454 .limit(
455 self.env
456 .opts
457 .hummock_time_travel_filter_out_objects_list_version_batch_size
458 as u64,
459 )
460 .stream(&self.env.meta_store_ref().conn)
461 .await?;
462 let mut next_prev_version_id = None;
463 while let Some(model) = version_stream.try_next().await? {
464 let version =
465 HummockVersion::from_persisted_protobuf_owned(model.version.to_protobuf());
466 for object_id in version.get_object_ids() {
467 result.remove(&object_id);
468 }
469 next_prev_version_id = Some(model.version_id);
470 }
471 if let Some(next_prev_version_id) = next_prev_version_id {
472 prev_version_id = Some(next_prev_version_id);
473 } else {
474 break;
475 }
476 }
477 }
478
479 {
481 let mut prev_version_id: Option<HummockVersionId> = None;
482 loop {
483 let query = hummock_time_travel_delta::Entity::find();
484 let query = if let Some(prev_version_id) = prev_version_id {
485 query.filter(hummock_time_travel_delta::Column::VersionId.gt(prev_version_id))
486 } else {
487 query
488 };
489 let mut version_stream = query
490 .order_by_asc(hummock_time_travel_delta::Column::VersionId)
491 .limit(
492 self.env
493 .opts
494 .hummock_time_travel_filter_out_objects_list_delta_batch_size
495 as u64,
496 )
497 .stream(&self.env.meta_store_ref().conn)
498 .await?;
499 let mut next_prev_version_id = None;
500 while let Some(model) = version_stream.try_next().await? {
501 let version_delta = HummockVersionDelta::from_persisted_protobuf_owned(
502 model.version_delta.to_protobuf(),
503 );
504 for object_id in version_delta.newly_added_object_ids(true) {
506 result.remove(&object_id);
507 }
508 next_prev_version_id = Some(model.version_id);
509 }
510 if let Some(next_prev_version_id) = next_prev_version_id {
511 prev_version_id = Some(next_prev_version_id);
512 } else {
513 break;
514 }
515 }
516 }
517
518 Ok(result)
519 }
520
521 pub(crate) async fn time_travel_pinned_object_count(&self) -> Result<u64> {
522 let count = hummock_sstable_info::Entity::find()
523 .count(&self.env.meta_store_ref().conn)
524 .await?;
525 Ok(count)
526 }
527
528 pub async fn epoch_to_version(
534 &self,
535 query_epoch: HummockEpoch,
536 table_id: TableId,
537 ) -> Result<HummockVersion> {
538 let sql_store = self.env.meta_store_ref();
539 let _permit = self.inflight_time_travel_query.try_acquire().map_err(|_| {
540 anyhow!(format!(
541 "too many inflight time travel queries, max_inflight_time_travel_query={}",
542 self.env.opts.max_inflight_time_travel_query
543 ))
544 })?;
545 let epoch_to_version = hummock_epoch_to_version::Entity::find()
546 .filter(
547 Condition::any()
548 .add(
549 hummock_epoch_to_version::Column::TableId
550 .eq(i64::from(table_id.as_raw_id())),
551 )
552 .add(hummock_epoch_to_version::Column::TableId.eq(0)),
554 )
555 .filter(
556 hummock_epoch_to_version::Column::Epoch
557 .lte(risingwave_meta_model::Epoch::try_from(query_epoch).unwrap()),
558 )
559 .order_by_desc(hummock_epoch_to_version::Column::Epoch)
560 .one(&sql_store.conn)
561 .await?
562 .ok_or_else(|| Error::TimeTravelVersionExpired {
563 table_id,
564 epoch: query_epoch,
565 })?;
566 let timer = self
567 .metrics
568 .time_travel_version_replay_latency
569 .start_timer();
570 let actual_version_id = epoch_to_version.version_id;
571 tracing::debug!(
572 query_epoch,
573 query_tz = ?(Epoch(query_epoch).as_timestamptz()),
574 actual_epoch = epoch_to_version.epoch,
575 actual_tz = ?(Epoch(u64::try_from(epoch_to_version.epoch).unwrap()).as_timestamptz()),
576 %actual_version_id,
577 "convert query epoch"
578 );
579
580 let Some(resolved) =
581 resolve_time_travel_version(&sql_store.conn, actual_version_id).await?
582 else {
583 return Err(Error::TimeTravelVersionExpired {
584 table_id,
585 epoch: query_epoch,
586 });
587 };
588 let mut actual_version = replay_archive(
589 resolved.replay_version.to_protobuf(),
590 resolved.deltas.into_iter().map(|delta| delta.to_protobuf()),
591 )?;
592 if actual_version.id != actual_version_id {
593 return Err(Error::TimeTravelVersionExpired {
594 table_id,
595 epoch: query_epoch,
596 });
597 }
598
599 let mut sst_ids = actual_version
601 .get_sst_ids()
602 .into_iter()
603 .collect::<VecDeque<_>>();
604 let sst_count = sst_ids.len();
605 let mut sst_id_to_info = HashMap::with_capacity(sst_count);
606 let sst_info_fetch_batch_size = self.env.opts.hummock_time_travel_sst_info_fetch_batch_size;
607 while !sst_ids.is_empty() {
608 let sst_infos = hummock_sstable_info::Entity::find()
609 .filter(hummock_sstable_info::Column::SstId.is_in(
610 sst_ids.drain(..std::cmp::min(sst_info_fetch_batch_size, sst_ids.len())),
611 ))
612 .all(&sql_store.conn)
613 .await?;
614 for sst_info in sst_infos {
615 let sst_info: SstableInfo = sst_info.sstable_info.to_protobuf().into();
616 sst_id_to_info.insert(sst_info.sst_id, sst_info);
617 }
618 }
619 if sst_count != sst_id_to_info.len() {
620 return Err(Error::TimeTravelVersionExpired {
621 table_id,
622 epoch: query_epoch,
623 });
624 }
625 refill_version(&mut actual_version, &sst_id_to_info, table_id);
626 timer.observe_duration();
627 Ok(actual_version)
628 }
629
630 pub(crate) async fn write_time_travel_metadata(
631 &self,
632 txn: &DatabaseTransaction,
633 version: Option<&HummockVersion>,
634 delta: HummockVersionDelta,
635 time_travel_table_ids: HashSet<StateTableId>,
636 skip_sst_ids: &HashSet<HummockSstableId>,
637 tables_to_commit: impl Iterator<Item = (&TableId, &CompactionGroupId, u64)>,
638 ) -> Result<Option<HashSet<HummockSstableId>>> {
639 let _timer = self
640 .metrics
641 .time_travel_write_metadata_latency
642 .start_timer();
643 if self
644 .env
645 .system_params_reader()
646 .await
647 .time_travel_retention_ms()
648 == 0
649 {
650 return Ok(None);
651 }
652 async fn write_sstable_infos(
653 mut sst_infos: impl Iterator<Item = &SstableInfo>,
654 txn: &DatabaseTransaction,
655 batch_size: usize,
656 ) -> Result<usize> {
657 let mut count = 0;
658 let mut is_finished = false;
659 while !is_finished {
660 let mut remain = batch_size;
661 let mut batch = vec![];
662 while remain > 0 {
663 let Some(sst_info) = sst_infos.next() else {
664 is_finished = true;
665 break;
666 };
667 batch.push(hummock_sstable_info::ActiveModel {
668 sst_id: Set(sst_info.sst_id),
669 object_id: Set(sst_info.object_id),
670 sstable_info: Set(SstableInfoV2Backend::from(&sst_info.to_protobuf())),
671 });
672 remain -= 1;
673 count += 1;
674 }
675 if batch.is_empty() {
676 break;
677 }
678 hummock_sstable_info::Entity::insert_many(batch)
679 .on_conflict_do_nothing()
680 .exec(txn)
681 .await?;
682 }
683 Ok(count)
684 }
685
686 let mut batch = vec![];
687 for (table_id, _cg_id, committed_epoch) in tables_to_commit {
688 if !time_travel_table_ids.contains(table_id) {
689 continue;
690 }
691 let m = hummock_epoch_to_version::ActiveModel {
692 epoch: Set(committed_epoch.try_into().unwrap()),
693 table_id: Set(i64::from(table_id.as_raw_id())),
694 version_id: Set(delta.id),
695 };
696 batch.push(m);
697 if batch.len()
698 >= self
699 .env
700 .opts
701 .hummock_time_travel_epoch_version_insert_batch_size
702 {
703 hummock_epoch_to_version::Entity::insert_many(std::mem::take(&mut batch))
705 .do_nothing()
706 .exec(txn)
707 .await?;
708 }
709 }
710 if !batch.is_empty() {
711 hummock_epoch_to_version::Entity::insert_many(batch)
713 .do_nothing()
714 .exec(txn)
715 .await?;
716 }
717
718 let mut version_sst_ids = None;
719 if let Some(version) = version {
720 version_sst_ids = Some(
722 version
723 .get_sst_infos()
724 .filter_map(|s| {
725 if s.table_ids
726 .iter()
727 .any(|tid| time_travel_table_ids.contains(tid))
728 {
729 return Some(s.sst_id);
730 }
731 None
732 })
733 .collect(),
734 );
735 write_sstable_infos(
736 version.get_sst_infos().filter(|s| {
737 !skip_sst_ids.contains(&s.sst_id)
738 && s.table_ids
739 .iter()
740 .any(|tid| time_travel_table_ids.contains(tid))
741 }),
742 txn,
743 self.env.opts.hummock_time_travel_sst_info_insert_batch_size,
744 )
745 .await?;
746 let m = hummock_time_travel_version::ActiveModel {
747 version_id: Set(version.id),
748 version: Set(
749 (&IncompleteHummockVersion::from((version, &time_travel_table_ids))
750 .to_protobuf())
751 .into(),
752 ),
753 };
754 hummock_time_travel_version::Entity::insert(m)
755 .on_conflict_do_nothing()
756 .exec(txn)
757 .await?;
758 return Ok(version_sst_ids);
760 }
761 let written = write_sstable_infos(
762 delta.newly_added_sst_infos(true).filter(|s| {
763 !skip_sst_ids.contains(&s.sst_id)
764 && s.table_ids
765 .iter()
766 .any(|tid| time_travel_table_ids.contains(tid))
767 }),
768 txn,
769 self.env.opts.hummock_time_travel_sst_info_insert_batch_size,
770 )
771 .await?;
772 let has_state_table_info_delta = delta
773 .state_table_info_delta
774 .keys()
775 .any(|table_id| time_travel_table_ids.contains(table_id));
776 if written > 0 || has_state_table_info_delta {
777 let m = hummock_time_travel_delta::ActiveModel {
778 version_id: Set(delta.id),
779 version_delta: Set((&IncompleteHummockVersionDelta::from((
780 &delta,
781 &time_travel_table_ids,
782 ))
783 .to_protobuf())
784 .into()),
785 };
786 hummock_time_travel_delta::Entity::insert(m)
787 .on_conflict_do_nothing()
788 .exec(txn)
789 .await?;
790 }
791
792 Ok(version_sst_ids)
793 }
794}
795
796struct ResolvedTimeTravelVersion {
797 replay_version: IncompleteHummockVersion,
798 deltas: Vec<IncompleteHummockVersionDelta>,
799}
800
801async fn resolve_time_travel_version(
802 conn: &impl ConnectionTrait,
803 actual_version_id: HummockVersionId,
804) -> Result<Option<ResolvedTimeTravelVersion>> {
805 let Some(replay_version) = hummock_time_travel_version::Entity::find()
806 .filter(hummock_time_travel_version::Column::VersionId.lte(actual_version_id))
807 .order_by_desc(hummock_time_travel_version::Column::VersionId)
808 .one(conn)
809 .await?
810 else {
811 return Ok(None);
812 };
813 let deltas = hummock_time_travel_delta::Entity::find()
814 .filter(hummock_time_travel_delta::Column::VersionId.gt(replay_version.version_id))
815 .filter(hummock_time_travel_delta::Column::VersionId.lte(actual_version_id))
816 .order_by_asc(hummock_time_travel_delta::Column::VersionId)
817 .all(conn)
818 .await?;
819 Ok(Some(ResolvedTimeTravelVersion {
820 replay_version: IncompleteHummockVersion::from_persisted_protobuf_owned(
821 replay_version.version.to_protobuf(),
822 ),
823 deltas: deltas
824 .into_iter()
825 .map(|delta| {
826 IncompleteHummockVersionDelta::from_persisted_protobuf_owned(
827 delta.version_delta.to_protobuf(),
828 )
829 })
830 .collect(),
831 }))
832}
833
834fn replay_archive(
836 version: PbHummockVersion,
837 deltas: impl Iterator<Item = PbHummockVersionDelta>,
838) -> Result<HummockVersion> {
839 let mut last_version = HummockVersion::from_persisted_protobuf_owned(version);
842 for d in deltas {
843 let d = HummockVersionDelta::from_persisted_protobuf_owned(d);
844 debug_assert!(
845 !should_mark_next_time_travel_version_snapshot(&d),
846 "unexpected time travel delta {:?}",
847 d
848 );
849 if d.prev_id < last_version.id {
850 return Err(Error::TimeTravel(anyhow!(format!(
851 "invalid time travel delta chain: delta {} has prev version {}, but replay has reached {}",
852 d.id, d.prev_id, last_version.id
853 ))));
854 }
855 last_version.id = d.prev_id;
858 last_version.apply_version_delta(&d);
859 }
860 Ok(last_version)
861}
862
863pub fn require_sql_meta_store_err() -> Error {
864 Error::TimeTravel(anyhow!("require SQL meta store"))
865}
866
867pub fn should_mark_next_time_travel_version_snapshot(delta: &HummockVersionDelta) -> bool {
869 delta.group_deltas.iter().any(|(_, deltas)| {
870 deltas
871 .group_deltas
872 .iter()
873 .any(|d| !matches!(d, GroupDeltaCommon::NewL0SubLevel(_)))
874 })
875}
876
877#[cfg(test)]
878mod tests {
879 use super::*;
880
881 fn version(id: u64) -> PbHummockVersion {
882 let mut version = HummockVersion::default();
883 version.id = id.into();
884 version.to_protobuf()
885 }
886
887 fn delta(id: u64, prev_id: u64) -> PbHummockVersionDelta {
888 let mut delta = HummockVersionDelta::default();
889 delta.id = id.into();
890 delta.prev_id = prev_id.into();
891 delta.to_protobuf()
892 }
893
894 #[test]
895 fn test_replay_archive_delta_chain() {
896 let replayed = replay_archive(version(1), [delta(4, 3), delta(5, 4)].into_iter()).unwrap();
897 assert_eq!(replayed.id, HummockVersionId::new(5));
898
899 assert!(replay_archive(version(3), [delta(4, 2)].into_iter()).is_err());
900 }
901}