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